← Files ModRetro Chromatic PluginARCHIVED FILE

scripts/device-mcp.mjs

17.6 KB · Oct 2, 2026 · 00:37 UTC

↓ Download file

#!/usr/bin/env node
/** Developer client: one foreground session, only the package's public MCP. */
import { createRequire } from "node:module";
import { readFile, realpath } from "node:fs/promises";
import path from "node:path";
import { createInterface } from "node:readline";
import { setTimeout as delay } from "node:timers/promises";
import { pathToFileURL } from "node:url";

const publicTools = ["device", "setup", "flash", "play", "rom_inspect"];
const isObject = (value) => value !== null && typeof value === "object" && !Array.isArray(value);
const isReadOnly = (tool, input) => tool === "rom_inspect" || (tool === "device" && ["status", "operation_status"].includes(input.command));
const message = (error) => error instanceof Error ? error.message : String(error);
const rejectedInput = (reason) => Object.assign(new Error(`Public tool input rejected before dispatch: ${reason}`), { code: "MCP_INPUT_INVALID", completion: "not-dispatched" });

// The published flat schemas keep tool clients compatible. Apply command-specific
// constraints before assuming ownership of a call the server cannot admit.
function validateCommand(tool, input) {
  if (tool === "device") {
    const allowed = { status: [], list_devices: ["requestId"], operation_status: ["operationId", "requestId"] };
    if (!Object.hasOwn(allowed, input.command)) throw rejectedInput("unknown device command");
    if (Object.keys(input).some((key) => key !== "command" && !allowed[input.command].includes(key))) throw rejectedInput("fields do not match the selected device command");
    if (input.command === "operation_status" && (input.operationId !== undefined) === (input.requestId !== undefined)) throw rejectedInput("operation_status needs exactly one of operationId or requestId");
  }
  if (tool === "setup") {
    const required = { install_drivers: "sessionId", detect_cartridge: "deviceToken" };
    if (!Object.hasOwn(required, input.command)) throw rejectedInput("unknown setup command");
    const field = required[input.command];
    if (input[field] === undefined || Object.keys(input).some((key) => !["command", "requestId", "confirm", field].includes(key))) throw rejectedInput(`the selected setup command requires only ${field} in addition to requestId and confirm`);
  }
  if (["setup", "flash", "play"].includes(tool) && input.confirm !== true) throw rejectedInput("the requested setup, live-play, or write action requires confirm:true");
}

/** Internal options supply test signals/timeouts; public calls retain their exact arguments. */
export async function openDeviceMcpSession({ packageRoot, workspaceRoot, signalSource = process, requestTimeoutMs = 60_000, maxBufferSize }) {
  const root = await realpath(packageRoot);
  const config = JSON.parse(await readFile(path.join(root, ".mcp.json"), "utf8")).mcpServers?.["modretro-chromatic"];
  if (!config || config.type !== "stdio" || typeof config.command !== "string"
    || !Array.isArray(config.args) || config.args.some((arg) => typeof arg !== "string")) {
    throw new Error("Select a compiled plugin with its ordinary modretro-chromatic stdio configuration.");
  }
  const runtime = await realpath(config.env?.GB_STUDIO_RUNTIME_ROOT ?? root);
  const manifest = JSON.parse(await readFile(path.join(runtime, "package.json"), "utf8"));
  let dependencies;
  if (manifest.bundledMcp === true) {
    dependencies = await import(pathToFileURL(path.join(runtime, "scripts/device-client-deps.mjs")).href);
  } else {
    const require = createRequire(path.join(runtime, "package.json"));
    dependencies = {
      ...await import(pathToFileURL(require.resolve("@modelcontextprotocol/sdk/client/index.js")).href),
      ...await import(pathToFileURL(require.resolve("@modelcontextprotocol/sdk/client/stdio.js")).href),
      ...await import(pathToFileURL(require.resolve("@modelcontextprotocol/sdk/validation/ajv")).href),
    };
  }
  const { Client, StdioClientTransport, getDefaultEnvironment, AjvJsonSchemaValidator } = dependencies;
  let pending; let inFlight = false; let stopped = false; let transportClosed = false;
  let closePromise; let deferredCloseCount = 0;
  const interruptions = []; const transportErrors = [];
  let interruptionCount = 0; let transportErrorCount = 0;
  // Also fences the SDK's own close-on-buffer-error path. No SDK internals,
  // detach, replacement server, or signal to an uncertain device operation.
  class DeviceTransport extends StdioClientTransport {
    async close() {
      if (pending && !transportClosed) { deferredCloseCount++; return; }
      return super.close();
    }
  }
  const transport = new DeviceTransport({
    command: config.command === "node" ? process.execPath : config.command,
    args: config.args, cwd: path.resolve(root, config.cwd ?? "."),
    env: { ...getDefaultEnvironment(), ...config.env,
      ...(workspaceRoot ? { GB_STUDIO_WORKSPACE_ROOT: await realpath(workspaceRoot) } : {}),
    }, stderr: "pipe", ...(maxBufferSize === undefined ? {} : { maxBufferSize }),
  });
  const stderrChunks = []; let stderrBytes = 0; let stderrRetainedBytes = 0;
  transport.stderr?.on("data", (chunk) => {
    const bytes = Buffer.from(chunk); stderrBytes += bytes.length;
    const keep = Math.min(bytes.length, 512 * 1024 - stderrRetainedBytes);
    if (keep > 0) { stderrChunks.push(Buffer.from(bytes.subarray(0, keep))); stderrRetainedBytes += keep; }
  });
  const signals = ["SIGINT", "SIGTERM", "SIGHUP"].map((signal) => [signal, () => {
    stopped = true; interruptionCount++;
    if (interruptions.length < 16) interruptions.push(signal);
    process.stderr.write("No new action will start. Waiting for the original public MCP operation; no retry or device reset.\n");
  }]);
  const removeSignals = () => { for (const [signal, handler] of signals) signalSource.off(signal, handler); };
  for (const [signal, handler] of signals) signalSource.on(signal, handler);
  transport.onclose = () => { transportClosed = true; removeSignals(); };
  transport.onerror = (error) => { transportErrorCount++; if (transportErrors.length < 16) transportErrors.push(message(error)); };
  const client = new Client({ name: "modretro-public-device-client", version: "0.1.0" });
  function observe(result, initial = false) {
    if (!pending) return;
    const value = result.structuredContent;
    if (!isObject(value)) return;
    if (initial && result.isError && value.operationStarted === false) { pending = undefined; return; }
    if (initial && typeof value.operationId === "string") pending.operationId = value.operationId;
    const matches = pending.operationId ? value.operationId === pending.operationId
      : pending.requestId && value.requestId === pending.requestId;
    if (matches && ["succeeded", "failed"].includes(value.state)) { pending = undefined; return; }
    if (matches && typeof value.operationId === "string") pending.operationId = value.operationId;
    // Recover only a unique new discovery in this same dedicated server.
    // The baseline distinguishes a completed discovery from old history;
    // neither idle status nor OPERATION_UNKNOWN clears an uncertain dispatch.
    const baseline = pending.discoveryBaseline;
    if (!pending.operationId && !pending.requestId && baseline && value.sessionId === baseline.sessionId && Array.isArray(value.operations)) {
      const candidates = value.operations.filter((operation) => isObject(operation) && operation.command === "list_devices"
        && typeof operation.operationId === "string" && !baseline.operationIds.includes(operation.operationId));
      if (candidates.length === 1) pending.operationId = candidates[0].operationId;
    }
  }
  const state = () => ({
    packageRoot: root, runtimeRoot: runtime, transportClosed, stopped, pendingOperation: pending ? { ...pending } : null,
    deferredCloseCount, stderr: Buffer.concat(stderrChunks).toString("utf8"), stderrBytes,
    stderrTruncated: stderrRetainedBytes < stderrBytes, interruptions: [...interruptions], interruptionCount,
    transportErrors: [...transportErrors], transportErrorCount,
  });
  async function close() {
    stopped = true;
    if (pending && !transportClosed) return { closeDeferred: true, completion: "unverified", ...state() };
    if (!closePromise) closePromise = (async () => {
      try { await client.close(); return { clientCloseResolved: true }; }
      catch (error) { return { clientCloseError: message(error) }; }
      finally { removeSignals(); }
    })();
    return { ...await closePromise, ...state() };
  }
  let tools; const inputValidators = new Map();
  try {
    await client.connect(transport);
    const listed = await client.listTools();
    tools = listed.tools.filter((entry) => ["device", "setup", "flash", "play"].includes(entry.name));
    if (tools.length !== 4 || listed.tools.some((entry) => entry.name.startsWith("chromatic_"))) throw new Error("The selected package does not expose device, setup, flash and play.");
    for (const entry of listed.tools.filter((entry) => publicTools.includes(entry.name))) {
      inputValidators.set(entry.name, new AjvJsonSchemaValidator().getValidator(entry.inputSchema));
    }
  } catch (error) {
    await close();
    throw error;
  }
  async function call(tool, input = {}, onAccepted) {
    if (!publicTools.includes(tool) || !isObject(input)) throw new Error("Use device, setup, flash, play, or rom_inspect with a JSON argument object.");
    if (inFlight) throw new Error("This client serializes requests in its original MCP session.");
    const validation = inputValidators.get(tool)?.(input);
    if (!validation?.valid) throw rejectedInput(validation?.errorMessage ?? "tool unavailable");
    validateCommand(tool, input);
    const readOnly = isReadOnly(tool, input);
    if ((stopped || pending) && !readOnly) throw new Error("A previous action is interrupted or unresolved. Only read its original status; do not start another action.");
    if (transportClosed) throw new Error("The original MCP transport closed. No replacement session or replay was started.");
    let initial; let result; let discoveryBaseline; let polls = 0; let dispatched = false;
    inFlight = true;
    try {
      let baseline;
      if (tool === "device" && input.command === "list_devices" && input.requestId === undefined) {
        discoveryBaseline = await client.callTool({ name: "device", arguments: { command: "status" } }, undefined, { timeout: requestTimeoutMs });
        const value = discoveryBaseline.structuredContent;
        if (discoveryBaseline.isError || typeof value?.sessionId !== "string" || !Array.isArray(value.operations)
          || value.operations.some((operation) => !isObject(operation) || typeof operation.operationId !== "string")) throw new Error("Could not bind this session's discovery history; no discovery was dispatched.");
        baseline = { sessionId: value.sessionId, operationIds: value.operations.map((operation) => operation.operationId) };
        if (stopped) throw new Error("Interrupted before discovery dispatch; no device action was started.");
      }
      if (!readOnly) pending = { tool, arguments: input, ...(baseline ? { discoveryBaseline: baseline } : {}),
        ...(typeof input.requestId === "string" ? { requestId: input.requestId } : {}) };
      dispatched = true;
      initial = await client.callTool({ name: tool, arguments: input }, undefined, { timeout: requestTimeoutMs });
      result = initial; observe(result, !readOnly);
      if (typeof result.structuredContent?.operationId === "string") onAccepted?.(result);
      while (result.structuredContent?.state === "running") {
        await delay(250);
        result = await client.callTool({ name: "device", arguments: { command: "operation_status", operationId: result.structuredContent.operationId } }, undefined, { timeout: requestTimeoutMs });
        polls++; observe(result);
      }
      return { tool, arguments: input, discoveryBaseline, initial, result, polls, completion: pending ? "unverified" : "observed", automaticRetry: false };
    } catch (error) {
      return { tool, arguments: input, discoveryBaseline, initial, result, polls, error: message(error),
        completion: dispatched ? "unverified" : "not-dispatched", automaticRetry: false };
    } finally { inFlight = false; }
  }
  async function reconcile() {
    if (!pending || transportClosed) return state();
    const query = pending.operationId ? { command: "operation_status", operationId: pending.operationId }
      : pending.requestId ? { command: "operation_status", requestId: pending.requestId } : { command: "status" };
    return call("device", query);
  }
  return { tools, call, close, reconcile, state };
}

/** One-shot mode is passive only; device workflows need the same-session API or --jsonl. */
export async function runDeviceMcp({ tool, arguments: input = {}, schemaOnly = false, ...options }) {
  if (!schemaOnly && (!isObject(input) || !isReadOnly(tool, input))) throw new Error("One-shot mode is passive only. Use --jsonl for discovery, setup, flash, and play in one MCP session.");
  const session = await openDeviceMcpSession(options);
  const result = schemaOnly ? { tools: session.tools } : await session.call(tool, input);
  return { ...result, ...await session.close() };
}

export async function runDeviceMcpJsonl(options, input = process.stdin, output = process.stdout) {
  const session = await openDeviceMcpSession(options);
  const lines = createInterface({ input, crlfDelay: Infinity, terminal: false });
  const iterator = lines[Symbol.asyncIterator]();
  const firstLine = iterator.next(); // Subscribe before ready can trigger a reply.
  const emit = (value) => output.write(JSON.stringify(value) + "\n");
  let pendingRequestId; let failed = false; let closing;
  emit({ event: "ready", ...session.state(), tools: session.tools });
  try {
    for (let item = await firstLine; !item.done; item = await iterator.next()) {
      const line = item.value;
      let request;
      try {
        if (Buffer.byteLength(line) > 65536) throw new Error("A client request must not exceed 64 KiB.");
        request = JSON.parse(line);
        if (!isObject(request) || typeof request.id !== "string" || request.id.length > 128) throw new Error("Each line needs a short string id.");
        if (request.close === true) break;
        if (!session.state().pendingOperation && isObject(request.arguments) && !isReadOnly(request.tool, request.arguments)) pendingRequestId = request.id;
        emit({ id: request.id, event: "dispatch", tool: request.tool, arguments: request.arguments });
        const result = await session.call(request.tool, request.arguments, (accepted) => emit({ id: request.id, event: "accepted", result: accepted }));
        if (result.error || result.result?.isError || result.completion === "unverified") failed = true;
        emit({ id: request.id, event: result.completion === "unverified" ? "unknown" : "result", ...result });
      } catch (error) { failed = true; emit({ ...(request?.id ? { id: request.id } : {}), event: "error", error: message(error),
        ...(error?.completion === "not-dispatched" ? { completion: "not-dispatched" } : {}), ...session.state() }); }
    }
  } finally {
    lines.close();
    closing = await session.close();
    if (closing.closeDeferred) emit({ event: "close-deferred", ...closing });
    // Remain a foreground owner after EOF: query only the original operation.
    // A broken transport/unknown result is never replaced by a forced close.
    let lastRecovery;
    while (closing.closeDeferred && !session.state().transportClosed) {
      await delay(1000);
      const original = session.state().pendingOperation;
      let recovered;
      try { recovered = await session.reconcile(); }
      catch (error) { recovered = { error: message(error), completion: "unverified" }; }
      const fingerprint = JSON.stringify(recovered);
      if (fingerprint !== lastRecovery) {
        emit({ id: pendingRequestId, event: recovered.completion === "observed" ? "result" : "unknown", recovered: true, original, ...recovered });
        lastRecovery = fingerprint;
      }
      closing = await session.close();
    }
    emit({ event: "closed", ...closing });
  }
  return { ...closing, failed: failed || Boolean(closing.pendingOperation || closing.clientCloseError) };
}

if (process.argv[1] && pathToFileURL(path.resolve(process.argv[1])).href === import.meta.url) {
  try {
    const args = process.argv.slice(2); const options = {};
    for (let index = 0; index < args.length; index++) {
      const key = args[index];
      if (key === "--schema") { options.schemaOnly = true; continue; }
      if (key === "--jsonl") { options.jsonl = true; continue; }
      if (!["--package", "--workspace", "--tool", "--arguments"].includes(key) || args[index + 1] === undefined) throw new Error("Usage: device-mcp.mjs --package <compiled-plugin> [--workspace <root>] (--jsonl | --schema | --tool <passive-tool> --arguments <JSON>)");
      const value = args[++index];
      if (key === "--package") options.packageRoot = value;
      if (key === "--workspace") options.workspaceRoot = value;
      if (key === "--tool") options.tool = value;
      if (key === "--arguments") options.arguments = JSON.parse(value);
    }
    if (!options.packageRoot) throw new Error("--package is required; no source or backend fallback is selected.");
    if (options.jsonl) {
      if (options.tool || options.schemaOnly) throw new Error("--jsonl owns one session; do not combine it with --tool or --schema.");
      const result = await runDeviceMcpJsonl(options);
      if (result.failed || result.interruptionCount) process.exitCode = 1;
    } else {
      const result = await runDeviceMcp(options);
      process.stdout.write(JSON.stringify(result) + "\n");
      if (result.error || result.result?.isError || result.result?.structuredContent?.success === false || result.interruptionCount) process.exitCode = 1;
    }
  } catch (error) { process.stderr.write(`${message(error)}\n`); process.exitCode = 1; }
}

SHA-256: e3c2f47a34563501c87bebd0f0d8df37a4aac5c732b135266919adae3a80c526