Developer SDK

JavaScript examples

Install the SDK · View source on GitHub ↗

messages/application.mjs

import {randomBytes} from "node:crypto";

function deferred() { let resolve; const promise = new Promise(done => {resolve = done;}); return {promise, resolve}; }
async function bounded(promise, millis, label) {
  let timer;
  try { return await Promise.race([promise, new Promise((_, reject) => {timer = setTimeout(() => reject(new Error(label + " timed out")), millis);})]); }
  finally { clearTimeout(timer); }
}
export function validateCommand(args) {
  if (!args.length || !["self", "peers", "watch", "send"].includes(args[0]) ||
      (args[0] === "send" ? args.length < 3 : args.length !== 1)) {
    throw new Error("usage: --self-test | --version | self | peers | watch | send CLIENT_ID TEXT");
  }
  if (args[0] === "send" && !/^[0-9a-f]{8}-(?:[0-9a-f]{4}-){3}[0-9a-f]{12}$/i.test(args[1])) throw new Error("CLIENT_ID must be a UUID");
}

// No socket emulation: the application handles discrete raw subprotocol frames.
// The same command controller is shared by JS and TS; each supplies its codec.
export async function runMessages(device, codec, args, options = {}) {
  validateCommand(args);
  const {log = console.log, error = console.error, timeoutMillis = 10000, signal} = options;
  const ready = deferred(), peersReady = deferred(), ack = deferred(), stopped = deferred();
  let subscription, pending, currentPeers, receivedPeers = false;
  let remoteConnected = device.getRemoteConnected();
  let lastPeers;
  const printPeers = peers => {
    currentPeers = remoteConnected ? peers : null;
    const value = !remoteConnected || peers == null ? "peers: unavailable" : "peers: " + JSON.stringify(peers);
    if (value !== lastPeers) {log(value); lastPeers = value;}
    if (remoteConnected && peers != null) {receivedPeers = true; peersReady.resolve();}
  };
  const remoteSub = device.addRemoteChangeListener(connected => {remoteConnected = connected; if (!connected) printPeers(null);});
  // Register live updates before the snapshot. The browser RPC resync also
  // replays the latest snapshot after registration reaches the companion.
  const peersSub = device.addNetworkPeersChangeListener(printPeers);
  const stop = () => stopped.resolve();
  signal?.addEventListener("abort", stop, {once: true});
  if (signal?.aborted) stop();
  try {
    const expires = Date.now() + 30000;
    while (!device.getRemoteConnected()) {
      if (device.getSyncError()) throw new Error(device.getSyncError());
      if (Date.now() >= expires) throw new Error("Native companion did not sync within 30 seconds; check its process, token, client JWT, and instance ID");
      if (signal?.aborted) return;
      await new Promise(resolve => setTimeout(resolve, 50));
    }
    remoteConnected = true;
    subscription = await device.enableSubprotocol(codec.SUBPROTOCOL, async message => {
      const owned = message.bytes.slice(); // Copy immediately before any await.
      await ready.promise;
      let frame;
      try {frame = codec.decode(owned);} catch (err) {error("rejected frame from " + message.sourceClientId + ": " + err.message); return;}
      if (frame.kind === 1) {
        log("text from " + message.sourceClientId + " id=" + frame.id + " " + JSON.stringify(frame.text));
        if (!await subscription.send(message.sourceClientId, codec.encode(2, frame.id))) throw new Error("SDK did not enqueue ACK");
      } else if (pending && message.sourceClientId.toLowerCase() === pending.clientId && frame.id === pending.id) {
        log("ACK from " + message.sourceClientId + " id=" + frame.id);
        ack.resolve();
      }
      // An ACK is never itself acknowledged.
    });
    ready.resolve();
    log("self: " + device.getClientId());
    if (args[0] === "self") return;
    printPeers(device.getNetworkPeers());
    const failed = subscription.closed.then(() => {throw new Error("Subprotocol subscription closed");});
    void failed.catch(() => {});
    if (args[0] === "peers") {
      if (!receivedPeers) await bounded(Promise.race([peersReady.promise, failed, stopped.promise]), timeoutMillis, "Peer snapshot");
    } else if (args[0] === "watch") {
      await Promise.race([failed, stopped.promise]);
    } else {
      const clientId = args[1].toLowerCase();
      if (!receivedPeers) await bounded(Promise.race([peersReady.promise, failed, stopped.promise]), timeoutMillis, "Peer snapshot");
      const connected = Array.isArray(currentPeers?.connected) ? currentPeers.connected : [];
      if (!connected.some(peer => typeof peer?.clientId === "string" && peer.clientId.toLowerCase() === clientId)) {
        throw new Error("Selected peer is not currently connected: " + args[1]);
      }
      const protocols = await Promise.race([subscription.querySubprotocols(clientId, timeoutMillis), failed]);
      if (protocols === null) throw new Error("Peer subprotocol query was unanswered or unavailable");
      if (!protocols.includes(codec.SUBPROTOCOL)) throw new Error("Peer does not advertise subprotocol " + codec.SUBPROTOCOL);
      let id; do {id = randomBytes(8).readBigUInt64BE();} while (id === 0n);
      const frame = codec.encode(1, id, args.slice(2).join(" "));
      pending = {clientId, id};
      if (!await subscription.send(clientId, frame)) throw new Error("SDK did not enqueue TEXT");
      log("sent id=" + id + "; waiting for ACK");
      await bounded(Promise.race([ack.promise, failed, stopped.promise]), timeoutMillis, "ACK");
    }
  } finally {
    pending = undefined;
    signal?.removeEventListener("abort", stop);
    try {
      await subscription?.close();
      await subscription?.closed.catch(() => {});
    } finally {peersSub(); remoteSub();}
  }
}