Developer SDK

JavaScript examples

Install the SDK · View source on GitHub ↗

socket/device.mjs

import {URNetwork} from "@urnetwork/sdk";
import {openDevice as openIntegratedDevice} from "../integration/client.mjs";
export function openDevice(config, wasmOptions = {}) {return openIntegratedDevice(URNetwork, config, wasmOptions);}

export async function directSocketEcho(device, protocol, endpoint, timeoutMillis = 30000) {
  if (protocol !== "tcp" && protocol !== "udp") throw new TypeError("Expected tcp or udp");
  const address = new URL(protocol + "://" + endpoint);
  if (!address.hostname || !address.port || address.username || address.password || address.search || address.hash || (address.pathname && address.pathname !== "/")) throw new TypeError("Expected host:port or [IPv6]:port");
  const host = address.hostname.replace(/^\[|\]$/g, ""), port = Number(address.port);
  const {TCPSocket, UDPSocket} = device.directSockets;
  const socket = protocol === "udp"
    ? new UDPSocket({remoteAddress: host, remotePort: port})
    : new TCPSocket(host, port);
  const {readable, writable} = await socket.opened;
  const reader = readable.getReader(), writer = writable.getWriter();
  let timer;
  try {
    const echo = async () => {
      const payload = new TextEncoder().encode("hello from URnetwork");
      await writer.write(protocol === "udp" ? {data: payload} : payload);
      if (protocol === "tcp") await writer.close();
      const parts = []; let received = 0;
      while (received < payload.length) {
        const result = await reader.read();
        if (result.done) throw new Error("Echo closed before the whole reply");
        const data = protocol === "udp" ? result.value.data : result.value;
        parts.push(data); received += data.length;
        if (protocol === "udp") break; // One UDPMessage is one datagram, even if empty.
      }
      const reply = new Uint8Array(received); let offset = 0;
      for (const part of parts) {reply.set(part, offset); offset += part.length;}
      return new TextDecoder().decode(reply);
    };
    return await Promise.race([echo(), new Promise((_, reject) => {
      timer = setTimeout(() => reject(new Error("Direct Sockets echo timed out")), timeoutMillis);
    })]);
  } finally {
    clearTimeout(timer);
    // Cancel/abort pending I/O before releasing the standard stream locks.
    await Promise.allSettled([reader.cancel(), writer.abort()]);
    reader.releaseLock(); writer.releaseLock();
    await socket.close(); await socket.closed;
  }
}