main
Colby Knox Publish 0.1.0 c1c3dea 26d ago
#!/usr/bin/env node
// MJPEG Stream Viewer: a raw multipart/x-mixed-replace camera stream opened
// straight in a Hydrogen webview paints at native size with scrollbars, and
// an http LAN camera cannot show inside https Hydrogen at all. This server
// proxies the stream (same origin, through the container), caps the frame
// rate, and serves a page that fits the picture to the pane.
//
// CLI (global --ai-thread first, like every Adom app):
//   mjpeg-stream-viewer --ai-thread "<name>" serve [--port N|auto] [--url <cam>] [--print-url] [--ttl H] [--keep-alive]
//   mjpeg-stream-viewer --ai-thread "<name>" show --url <cam> [--fps N] [--bare] [--name "<tab>"]
//   mjpeg-stream-viewer ls | state [--port N] | stop [--port N]
import http from "node:http";
import https from "node:https";
import fs from "node:fs";
import path from "node:path";
import os from "node:os";
import { execFile } from "node:child_process";
import { fileURLToPath } from "node:url";

const __dirname = path.dirname(fileURLToPath(import.meta.url));
const PKG = JSON.parse(fs.readFileSync(path.join(__dirname, "package.json"), "utf8"));
const VERSION = PKG.version;
const APP = "mjpeg-stream-viewer";
const HOME = os.homedir();
const INST_DIR = path.join(HOME, ".adom", "instances", APP);
const SETTINGS_FILE = path.join(INST_DIR, "settings.json");
const DEFAULT_PORT = 8873;

// ---- args -----------------------------------------------------------------
const argv = process.argv.slice(2);
let aiThread = process.env.ADOM_AI_THREAD || "";
if (argv[0] === "--ai-thread") { aiThread = argv[1] || ""; argv.splice(0, 2); }
const verb = argv.shift() || "serve";
const flags = {};
for (let i = 0; i < argv.length; i++) {
  const a = argv[i];
  if (!a.startsWith("--")) continue;
  const k = a.slice(2);
  if (["print-url", "keep-alive", "json", "bare"].includes(k)) { flags[k] = true; continue; }
  flags[k] = argv[i + 1]; i++;
}
const die = (msg, code = 1) => { console.error(msg); process.exit(code); };
const threadSlug = (t) => String(t || "anonymous").toLowerCase().replace(/[^a-z0-9]+/g, "-").replace(/^-|-$/g, "").slice(0, 60) || "anonymous";

// ---- settings (server side, per app-creator-settings) -----------------------
const DEFAULTS = { maxFps: 10, fit: "contain", ttlHours: 24, sources: [] };
function loadSettings() {
  try { return { ...DEFAULTS, ...JSON.parse(fs.readFileSync(SETTINGS_FILE, "utf8")) }; } catch { return { ...DEFAULTS }; }
}
function saveSettings(s) { fs.mkdirSync(INST_DIR, { recursive: true }); fs.writeFileSync(SETTINGS_FILE, JSON.stringify(s, null, 2)); }

// ---- registry (per app-creator-server-lifecycle) ----------------------------
function regFile(t) { return path.join(INST_DIR, threadSlug(t) + ".json"); }
function readRegistry() {
  try { return fs.readdirSync(INST_DIR).filter(f => f.endsWith(".json") && f !== "settings.json").map(f => { try { return JSON.parse(fs.readFileSync(path.join(INST_DIR, f), "utf8")); } catch { return null; } }).filter(Boolean); } catch { return []; }
}
function alive(pid) { try { process.kill(pid, 0); return true; } catch { return false; } }
function probe(port, ms = 800) {
  return new Promise((resolve) => {
    const req = http.get({ host: "127.0.0.1", port, path: "/version", timeout: ms }, (res) => { let b = ""; res.on("data", d => b += d); res.on("end", () => { try { resolve(JSON.parse(b)); } catch { resolve(null); } }); });
    req.on("error", () => resolve(null)); req.on("timeout", () => { req.destroy(); resolve(null); });
  });
}

// ---- URL policy --------------------------------------------------------------
function checkSourceUrl(u) {
  let url;
  try { url = new URL(String(u || "")); } catch { return { ok: false, reason: "not a URL" }; }
  if (url.protocol !== "http:" && url.protocol !== "https:") return { ok: false, reason: "only http and https sources" };
  return { ok: true, url };
}

// ---- the proxy: forward frames at a capped rate --------------------------------
// Upstream is a multipart stream of JPEG parts. Frames are found by the JPEG
// markers themselves (FFD8 .. FFD9), so it works with or without
// Content-Length headers and with any boundary string. Every forwarded frame is
// re-wrapped with our own boundary so the browser sees a clean stream.
const stats = { source: null, fps: 0, bytesPerSec: 0, width: 0, height: 0, clients: 0, lastFrameAt: 0, lastError: null };
let frameCount = 0, byteCount = 0;
setInterval(() => { stats.fps = frameCount; stats.bytesPerSec = byteCount; frameCount = 0; byteCount = 0; }, 1000).unref();

function jpegSize(buf) {
  let k = 2;
  while (k + 9 < buf.length) {
    if (buf[k] !== 0xFF) { k++; continue; }
    const m = buf[k + 1];
    if (m === 0xC0 || m === 0xC1 || m === 0xC2) return { height: buf.readUInt16BE(k + 5), width: buf.readUInt16BE(k + 7) };
    if (m === 0xD8 || (m >= 0xD0 && m <= 0xD7) || m === 0x01) { k += 2; continue; }
    k += 2 + buf.readUInt16BE(k + 2);
  }
  return null;
}

function proxyStream(sourceUrl, maxFps, res, onlyOne = false) {
  const lib = sourceUrl.protocol === "https:" ? https : http;
  const minGap = maxFps > 0 ? 1000 / maxFps : 0;
  const BOUNDARY = "mjpegviewerframe";
  let up = null, buf = Buffer.alloc(0), lastSent = 0, sentAny = false, closed = false;
  const end = () => { if (closed) return; closed = true; try { up && up.destroy(); } catch {} try { if (!res.writableEnded) res.end(); } catch {} stats.clients = Math.max(0, stats.clients - 1); };
  const req = lib.get(sourceUrl, { headers: { "User-Agent": `${APP}/${VERSION}` }, timeout: 8000 }, (u) => {
    up = u;
    if (u.statusCode !== 200) { stats.lastError = `upstream ${u.statusCode}`; res.writeHead(502, { "Content-Type": "text/plain" }); res.end(`upstream answered ${u.statusCode}`); return end(); }
    const ct = String(u.headers["content-type"] || "");
    stats.lastError = null;
    if (!onlyOne) res.writeHead(200, { "Content-Type": `multipart/x-mixed-replace; boundary=${BOUNDARY}`, "Cache-Control": "no-store", "Connection": "keep-alive" });
    stats.clients++;
    // A plain image (a camera that serves a still) is forwarded once and refreshed by the page.
    if (/^image\/jpeg/i.test(ct) && !/multipart/i.test(ct)) {
      const chunks = []; u.on("data", d => chunks.push(d)); u.on("end", () => { const f = Buffer.concat(chunks); emit(f); end(); }); return;
    }
    u.on("data", (d) => {
      buf = buf.length ? Buffer.concat([buf, d]) : d;
      // pull every complete JPEG out of the buffer
      for (;;) {
        const s = buf.indexOf(Buffer.from([0xFF, 0xD8]));
        if (s < 0) { if (buf.length > 4e6) buf = Buffer.alloc(0); break; }
        const e = buf.indexOf(Buffer.from([0xFF, 0xD9]), s + 2);
        if (e < 0) { if (s > 0) buf = buf.subarray(s); if (buf.length > 12e6) buf = Buffer.alloc(0); break; }
        const frame = buf.subarray(s, e + 2); buf = buf.subarray(e + 2);
        emit(frame);
        if (onlyOne) return end();
      }
    });
    u.on("end", end); u.on("error", (e) => { stats.lastError = e.message; end(); });
  });
  req.on("error", (e) => { stats.lastError = e.message; if (!res.headersSent) { res.writeHead(502, { "Content-Type": "text/plain" }); } res.end(`cannot reach ${sourceUrl.href}: ${e.message}`); end(); });
  req.on("timeout", () => { req.destroy(new Error("timeout")); });
  res.on("close", end);
  function emit(frame) {
    const now = Date.now();
    if (!onlyOne && sentAny && minGap && now - lastSent < minGap) return; // decimate
    lastSent = now; sentAny = true; frameCount++; byteCount += frame.length; stats.lastFrameAt = now;
    const sz = jpegSize(frame); if (sz) { stats.width = sz.width; stats.height = sz.height; }
    if (onlyOne) { res.writeHead(200, { "Content-Type": "image/jpeg", "Content-Length": frame.length, "Cache-Control": "no-store" }); res.end(frame); return; }
    res.write(`--${BOUNDARY}\r\nContent-Type: image/jpeg\r\nContent-Length: ${frame.length}\r\n\r\n`); res.write(frame); res.write("\r\n");
  }
}

// ---- server ----------------------------------------------------------------
async function serve() {
  const settings = loadSettings();
  const ttlH = Number(flags.ttl || settings.ttlHours || 24);
  const keepAlive = !!flags["keep-alive"];
  const initialUrl = flags.url || null;
  // reuse instead of respawn: this thread's live instance answers already
  for (const r of readRegistry()) {
    if (r.aiThread === aiThread && alive(r.pid) && await probe(r.port)) {
      const out = { ok: true, reused: true, port: r.port, url: `http://127.0.0.1:${r.port}/`, pid: r.pid, version: r.version };
      if (initialUrl) out.url += `?url=${encodeURIComponent(initialUrl)}`;
      if (flags["print-url"]) console.log(out.url); else console.log(JSON.stringify(out, null, 2));
      return;
    }
  }
  let port = flags.port === "auto" ? 0 : Number(flags.port || DEFAULT_PORT);
  const pageFile = path.join(__dirname, "public", "index.html");
  const icon = fs.readFileSync(path.join(__dirname, "public", "icon.svg"));
  let lastActivity = Date.now();
  const json = (res, code, obj) => { res.writeHead(code, { "Content-Type": "application/json", "Cache-Control": "no-store" }); res.end(JSON.stringify(obj)); };
  const readBody = (req) => new Promise((resolve) => { let b = ""; req.on("data", d => { b += d; if (b.length > 1e6) req.destroy(); }); req.on("end", () => { try { resolve(JSON.parse(b || "{}")); } catch { resolve({}); } }); });
  const needThread = (body, res) => {
    if (body && body.aiThread) return true;
    json(res, 400, { ok: false, error: "caller_identity_required", hint: `This call changes state, so it needs the ai-thread name. Re-run with: ${APP} --ai-thread "<name>" ...  (raw callers: add "aiThread" to the JSON body).` });
    return false;
  };
  const server = http.createServer(async (req, res) => {
    lastActivity = Date.now();
    const url = new URL(req.url, "http://x");
    const p = url.pathname.replace(/\/+$/, "") || "/";
    if (p === "/" || p === "/index.html") { res.writeHead(200, { "Content-Type": "text/html; charset=utf-8", "Cache-Control": "no-store" }); return res.end(fs.readFileSync(pageFile)); }
    if (p === "/icon.svg") { res.writeHead(200, { "Content-Type": "image/svg+xml", "Cache-Control": "public, max-age=3600" }); return res.end(icon); }
    if (p === "/version" || p === "/health") return json(res, 200, { ok: true, app: APP, version: VERSION, port, aiThread, pid: process.pid });
    if (p === "/settings" && req.method === "GET") return json(res, 200, { ...loadSettings(), _enums: { fit: ["contain", "cover"] }, _ranges: { maxFps: [1, 30], ttlHours: [1, 720] } });
    if (p === "/settings" && req.method === "POST") {
      const body = await readBody(req); if (!needThread(body, res)) return;
      const s = loadSettings();
      if (body.maxFps !== undefined) s.maxFps = Math.min(30, Math.max(1, Number(body.maxFps) || 10));
      if (body.fit !== undefined) s.fit = body.fit === "cover" ? "cover" : "contain";
      if (body.ttlHours !== undefined) s.ttlHours = Math.min(720, Math.max(1, Number(body.ttlHours) || 24));
      if (Array.isArray(body.sources)) s.sources = body.sources.filter(x => x && typeof x.url === "string" && checkSourceUrl(x.url).ok).map(x => ({ name: String(x.name || x.url).slice(0, 60), url: x.url })).slice(0, 32);
      saveSettings(s); return json(res, 200, { ok: true, ...s });
    }
    if (p === "/api/status") {
      const s = loadSettings();
      const fresh = stats.lastFrameAt && Date.now() - stats.lastFrameAt < 5000;
      return json(res, 200, {
        ok: true, version: VERSION, aiThread, state: stats.lastError ? "error" : (fresh ? "ok" : (stats.clients ? "connecting" : "idle")),
        leds: {
          source: { state: stats.source ? (stats.lastError ? "red" : (fresh ? "green" : "yellow")) : "grey", label: stats.source || "no source", _hint: stats.lastError || (stats.source ? "frames flowing" : "open ?url=<camera stream> or pick a saved source") },
          frames: { state: fresh ? "green" : "grey", label: `${stats.fps} fps, ${(stats.bytesPerSec / 1024).toFixed(0)} KB/s`, _hint: `capped at ${s.maxFps} fps by the proxy; raise or lower it in Settings` },
          size: { state: stats.width ? "green" : "grey", label: stats.width ? `${stats.width}x${stats.height}` : "no frame yet", _hint: stats.width > 1920 ? "a 4K source costs bandwidth for a pane this size; a 1080p source would look the same here" : "source frame size" },
        },
        stats: { ...stats }, settings: { maxFps: s.maxFps, fit: s.fit },
      });
    }
    if (p === "/stream" || p === "/snapshot.jpg") {
      const chk = checkSourceUrl(url.searchParams.get("url"));
      if (!chk.ok) { res.writeHead(400, { "Content-Type": "text/plain" }); return res.end(`bad url: ${chk.reason}`); }
      if (chk.url.hostname === "127.0.0.1" && Number(chk.url.port) === port) { res.writeHead(400, { "Content-Type": "text/plain" }); return res.end("that is this viewer"); }
      stats.source = chk.url.href;
      const fps = Math.min(30, Math.max(0, Number(url.searchParams.get("fps") || loadSettings().maxFps || 10)));
      return proxyStream(chk.url, fps, res, p === "/snapshot.jpg");
    }
    if (p === "/api/snapshot" && req.method === "POST") {
      const body = await readBody(req); if (!needThread(body, res)) return;
      const chk = checkSourceUrl(body.url); if (!chk.ok) return json(res, 400, { ok: false, error: chk.reason });
      const dir = path.join(HOME, "project", "screenshots"); fs.mkdirSync(dir, { recursive: true });
      const file = path.join(dir, `mjpeg-${new Date().toISOString().replace(/[:.]/g, "-")}.jpg`);
      const chunks = [];
      // A stand-in response that collects one frame; end() answers exactly
      // once (the proxy calls end on both the frame and the upstream close).
      let answered = false;
      const fake = { headersSent: false, writeHead() {}, write(d) { chunks.push(Buffer.from(d)); }, on() {},
        end(d) {
          if (answered) return; answered = true;
          if (d) chunks.push(Buffer.from(d)); const b = Buffer.concat(chunks);
          if (b.length > 100) { fs.writeFileSync(file, b); json(res, 200, { ok: true, file, bytes: b.length, hint: "Read the file to view it; it is under ~/project/screenshots so VS Code shows it too." }); }
          else json(res, 502, { ok: false, error: "no frame received", hint: stats.lastError || "is the camera reachable from this container?" });
        } };
      return proxyStream(chk.url, 0, fake, true);
    }
    if (p === "/shutdown" && req.method === "POST") {
      const body = await readBody(req); if (!needThread(body, res)) return;
      json(res, 200, { ok: true, bye: true }); setTimeout(() => { deregister(); process.exit(0); }, 100); return;
    }
    res.writeHead(404, { "Content-Type": "text/plain" }); res.end("not found");
  });
  const deregister = () => { try { fs.unlinkSync(regFile(aiThread)); } catch {} };
  server.listen(port, "127.0.0.1", () => {
    port = server.address().port;
    fs.mkdirSync(INST_DIR, { recursive: true });
    fs.writeFileSync(regFile(aiThread), JSON.stringify({ app: APP, aiThread, port, pid: process.pid, version: VERSION, startedAt: new Date().toISOString(), cwd: process.cwd() }, null, 2));
    let url = `http://127.0.0.1:${port}/`; if (initialUrl) url += `?url=${encodeURIComponent(initialUrl)}`;
    if (flags["print-url"]) console.log(url); else console.log(JSON.stringify({ ok: true, port, url, pid: process.pid, aiThread, ttlHours: keepAlive ? null : ttlH }, null, 2));
    // the dock pre-opens a tab named by webview.title; claim it once we are up
    if (process.env.ADOM_DOCK_WEBVIEW) execFile("adom-cli", ["hydrogen", "webview", "open-or-refresh", "--name", process.env.ADOM_DOCK_WEBVIEW, "--url", url], () => {});
    if (!keepAlive) setInterval(() => { if (Date.now() - lastActivity > ttlH * 3600e3) { deregister(); process.exit(0); } }, 60e3).unref();
  });
  process.on("SIGTERM", () => { deregister(); process.exit(0); }); process.on("SIGINT", () => { deregister(); process.exit(0); });
}

async function findInstance() {
  if (flags.port) return { port: Number(flags.port) };
  for (const r of readRegistry()) if (alive(r.pid) && await probe(r.port)) return r;
  return null;
}
async function main() {
  if (verb === "serve") { if (!aiThread) die(`caller_identity_required: serve needs the ai-thread name.\n  Re-run: ${APP} --ai-thread "<name>" serve ...`); return serve(); }
  if (verb === "ls") {
    const rows = [];
    for (const r of readRegistry()) rows.push({ ...r, alive: alive(r.pid) && !!(await probe(r.port)), drift: r.version !== VERSION ? `running ${r.version}, installed ${VERSION}` : null });
    console.log(JSON.stringify({ ok: true, instances: rows }, null, 2)); return;
  }
  if (verb === "state") { const i = await findInstance(); if (!i) die("no running instance (start one: serve)"); const s = await new Promise((r) => http.get(`http://127.0.0.1:${i.port}/api/status`, (res) => { let b = ""; res.on("data", d => b += d); res.on("end", () => r(b)); })); console.log(s); return; }
  if (verb === "stop") {
    if (!aiThread) die(`caller_identity_required: stop needs the ai-thread name.`);
    const i = await findInstance(); if (!i) die("no running instance");
    const req = http.request({ host: "127.0.0.1", port: i.port, path: "/shutdown", method: "POST", headers: { "Content-Type": "application/json" } }, (res) => { let b = ""; res.on("data", d => b += d); res.on("end", () => console.log(b)); });
    req.end(JSON.stringify({ aiThread })); return;
  }
  if (verb === "show") {
    if (!aiThread) die(`caller_identity_required: show needs the ai-thread name.\n  Re-run: ${APP} --ai-thread "<name>" show --url <camera stream>`);
    const src = flags.url; if (!src || !checkSourceUrl(src).ok) die("show needs --url <http(s) camera stream>");
    let inst = await findInstance();
    if (!inst) {
      // start our own server detached and wait for it
      const child = (await import("node:child_process")).spawn(process.execPath, [fileURLToPath(import.meta.url), "--ai-thread", aiThread, "serve", "--port", "auto", "--print-url"], { detached: true, stdio: ["ignore", "pipe", "inherit"] });
      const line = await new Promise((r) => { let b = ""; child.stdout.on("data", d => { b += d; if (b.includes("\n")) r(b.trim()); }); setTimeout(() => r(b.trim()), 8000); });
      child.unref();
      const m = /127\.0\.0\.1:(\d+)/.exec(line); if (!m) die("could not start the viewer: " + line);
      inst = { port: Number(m[1]) };
    }
    const q = `?url=${encodeURIComponent(src)}${flags.fps ? `&fps=${encodeURIComponent(flags.fps)}` : ""}${flags.bare ? "&bare=1" : ""}`;
    const proxyBase = (process.env.VSCODE_PROXY_URI || "").replace("{{port}}", String(inst.port));
    const pageUrl = proxyBase ? proxyBase.replace(/\/?$/, "/") + q : `http://127.0.0.1:${inst.port}/${q}`;
    const name = flags.name || "MJPEG Stream Viewer";
    execFile("adom-cli", ["hydrogen", "webview", "open-or-refresh", "--name", name, "--url", pageUrl, "--display-icon", "mdi:cctv"], (err, out) => {
      console.log(JSON.stringify({ ok: !err, port: inst.port, url: pageUrl, tab: name, hint: err ? `open this URL in a Hydrogen webview tab (adom-workspace-control skill): ${pageUrl}` : "opened in a Hydrogen webview tab in a work pane (never the VS Code pane)" }, null, 2));
    });
    return;
  }
  die(`unknown verb '${verb}'. Verbs: serve, show, ls, state, stop. Global flag: --ai-thread "<name>" first.`, 2);
}
main().catch((e) => die(String(e && e.message || e)));