../

WebSockets

The WebSocket protocol on the wire, the browser WebSocket API, typed message protocols, and production servers on Bun and Node: heartbeats, backpressure, auth, scaling and security. Building a multiplayer app on top of it (Socket.IO, Durable Objects, chat, game worlds, collaboration) is under Realtime.

Protocol

A WebSocket starts life as an HTTP/1.1 request that asks to switch protocols (RFC 6455). After the 101 response the TCP connection carries framed, bidirectional messages with 2 to 14 bytes of overhead per frame and no further HTTP headers.

Handshake

GET /chat HTTP/1.1
Host: api.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Version: 13
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Origin: https://app.example.com
Sec-WebSocket-Protocol: chat.v2, chat.v1
Sec-WebSocket-Extensions: permessage-deflate
 
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
Sec-WebSocket-Protocol: chat.v2
HeaderSent byPurpose
Upgrade: websocketbothrequest / confirm the protocol switch
Connection: Upgradebothmarks Upgrade as hop-by-hop
Sec-WebSocket-Versionclientalways 13; server answers 426 with supported versions
Sec-WebSocket-Keyclient16 random bytes, base64; proves the server speaks WebSocket
Sec-WebSocket-Acceptserverbase64 of SHA-1 of key + fixed GUID
Originbrowserpage origin; the server must check it (no CORS here)
Sec-WebSocket-Protocolbothclient offers subprotocols, server picks at most one
Sec-WebSocket-Extensionsbothnegotiated extensions, e.g. permessage-deflate

The accept value is not security, only proof that the server understood the upgrade:

accept.ts
const GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11";
 
async function acceptFor(key: string): Promise<string> {
  const bytes = new TextEncoder().encode(key + GUID);
  const hash = await crypto.subtle.digest("SHA-1", bytes);
  return btoa(String.fromCharCode(...new Uint8Array(hash)));
}
 
// RFC 6455 example: resolves to
// "s3pPLMBiTxaQ9kYGzzhZRbK+xOo="
await acceptFor("dGhlIHNhbXBsZSBub25jZQ==");

Frames

FieldBitsMeaning
FIN1last frame of a message
RSV1–33extension flags (RSV1 = compressed by permessage-deflate)
opcode4frame type, see below
MASK1payload is masked; must be 1 client to server
Payload len7, +16, +640–125 literal; 126: next 16 bits; 127: next 64 bits
Masking key0 or 32present when MASK = 1
Payloadnextension data + application data
OpcodeFrameNotes
0x0continuationlater fragments of a fragmented message
0x1textpayload must be valid UTF-8 (else close 1007)
0x2binaryarbitrary bytes
0x3–0x7reservedfuture data frames
0x8closeoptional 2-byte code + UTF-8 reason
0x9pingpeer must answer with a pong carrying the same payload
0xApongreply to ping, or unsolicited one-way heartbeat
0xB–0xFreservedfuture control frames
  • Control frames (close, ping, pong) carry at most 125 bytes and are never fragmented; they may be interleaved between fragments of a data message.
  • Masking: the client XORs each payload byte with a fresh random 4-byte key (byte[i] ^ key[i % 4]). It stops attacker-chosen bytes from poisoning intermediary caches; it is not encryption. Servers never mask and must fail a connection that receives an unmasked client frame.
  • Messages are delivered whole to the application: the browser API has no access to frames, fragmentation or ping/pong.
  • Head-of-line blocking: one large message delays everything behind it on the same socket.

URLs, subprotocols, extensions

SchemePortTransportUse
ws://80plain TCPlocal development only
wss://443TLS over TCPalways in production; proxies pass it through unmodified
  • An https: page cannot open ws:// (mixed content). Relative http(s): URLs are accepted by the constructor and mapped to ws(s): in current browsers.
  • Subprotocols are application-level version tags: new WebSocket(url, ["chat.v2", "chat.v1"]). The server echoes one; ws.protocol is "" if none was chosen.
  • permessage-deflate (RFC 7692) compresses each message; trades CPU and memory per connection for bandwidth.
  • RFC 8441 (HTTP/2) and RFC 9220 (HTTP/3) can bootstrap WebSockets over one multiplexed connection with extended CONNECT; the API is unchanged.

Close codes

CodeNameMeaning
1000Normal closuredone, purpose fulfilled
1001Going awaypage navigated away, server shutting down
1002Protocol errormalformed frame, unmasked client frame
1003Unsupported datae.g. binary sent to a text-only endpoint
1005No status receivedreserved: close frame had no code (never sent)
1006Abnormal closurereserved: no close frame, e.g. TCP dropped
1007Invalid payloadtext frame was not UTF-8
1008Policy violationgeneric "not allowed" (auth failed, bad origin)
1009Message too bigexceeded the peer's size limit
1010Mandatory extensionclient needed an extension the server refused
1011Internal errorunexpected server condition
1012Service restartreconnect later (IANA registry)
1013Try again lateroverloaded or rate limited (IANA registry)
1015TLS handshake failurereserved: reported locally only
3000–3999registeredlibraries and frameworks, via IANA
4000–4999privateyour application's own codes

Browser API

client.ts
const ws = new WebSocket("wss://api.example.com/chat", [
  "chat.v2",
]);
ws.binaryType = "arraybuffer"; // default "blob"
 
ws.addEventListener("open", () => {
  ws.send(JSON.stringify({ type: "join", room: "lobby" }));
});
 
ws.addEventListener("message", (e: MessageEvent) => {
  if (typeof e.data === "string") handleText(e.data);
  else handleBinary(new Uint8Array(e.data as ArrayBuffer));
});
 
ws.addEventListener("error", () => {
  // no details by design; a close event always follows
});
 
ws.addEventListener("close", (e) => {
  console.info(e.code, e.reason, e.wasClean);
});
 
// later
ws.close(1000, "bye");

readyState

ValueConstantMeaningsend()
0WebSocket.CONNECTINGhandshake in progressthrows InvalidStateError
1WebSocket.OPENreadyqueues the data
2WebSocket.CLOSINGclose handshake startedsilently discarded
3WebSocket.CLOSEDclosed or failed to opensilently discarded

Members

MemberTypeNotes
new WebSocket(url, protocols?)WebSocketconnects immediately; no custom headers possible
send(data)voidstring, ArrayBuffer, Blob, TypedArray, DataView
close(code?, reason?)voidcode 1000 or 3000–4999; reason max 123 UTF-8 bytes
binaryType"blob" | "arraybuffer"type of binary e.data
bufferedAmountnumberbytes queued by send() but not yet sent
readyState0 | 1 | 2 | 3see above
protocol, extensionsstringwhat the server accepted
urlstringresolved absolute URL
onopen … onclosehandler propsor addEventListener
EventObjectUseful fields
openEventnone
messageMessageEventdata: string | Blob | ArrayBuffer, origin
errorEventnone; deliberately opaque
closeCloseEventcode, reason, wasClean
  • Calling close() with any other code throws InvalidAccessError; a longer reason throws SyntaxError.
  • JS cannot send ping frames; browsers answer server pings automatically.
  • There is no auto-reconnect and no delivery acknowledgement: build both yourself (see below).

Typed messages

Define one discriminated union per direction, validate everything that arrives, and funnel every outgoing message through a typed helper.

protocol.ts
import { z } from "zod";
 
export const ClientMsg = z.discriminatedUnion("type", [
  z.object({
    type: z.literal("join"),
    room: z.string().min(1).max(64),
  }),
  z.object({
    type: z.literal("chat"),
    room: z.string().min(1).max(64),
    text: z.string().min(1).max(2000),
  }),
  z.object({ type: z.literal("ping"), t: z.number() }),
]);
export type ClientMsg = z.infer<typeof ClientMsg>;
 
export type ServerMsg =
  | { type: "joined"; room: string; members: number }
  | {
      type: "chat";
      room: string;
      from: string;
      text: string;
    }
  | { type: "pong"; t: number }
  | { type: "error"; code: "bad_message" | "forbidden" };
 
export function parseClientMsg(
  raw: unknown,
): ClientMsg | null {
  if (typeof raw !== "string") return null;
  try {
    const result = ClientMsg.safeParse(JSON.parse(raw));
    return result.success ? result.data : null;
  } catch {
    return null; // not JSON
  }
}
 
interface Sink {
  send(data: string): unknown;
}
 
export function sendMsg(ws: Sink, msg: ServerMsg): void {
  ws.send(JSON.stringify(msg));
}

Exhaustive dispatch: adding a new variant to ClientMsg breaks the build until it is handled.

dispatch.ts
import {
  type ClientMsg,
  type ServerMsg,
  parseClientMsg,
  sendMsg,
} from "./protocol.js";
 
type Conn = { send(data: string): unknown; userId: string };
 
export function onMessage(conn: Conn, raw: unknown): void {
  const msg = parseClientMsg(raw);
  if (!msg) {
    return sendMsg(conn, {
      type: "error",
      code: "bad_message",
    });
  }
  switch (msg.type) {
    case "join":
      return sendMsg(conn, joinRoom(conn, msg.room));
    case "chat":
      return broadcast(msg.room, {
        type: "chat",
        room: msg.room,
        from: conn.userId,
        text: msg.text,
      });
    case "ping":
      return sendMsg(conn, { type: "pong", t: msg.t });
    default:
      msg satisfies never;
  }
}
  • Share protocol.ts between client and server so both sides agree on the shape.
  • Add a version field or subprotocol (chat.v2) before the first breaking change, not after.
  • For binary or high-volume protocols, swap JSON for MessagePack or Protocol Buffers; keep the same union-of-messages design.
  • Correlate request/response pairs with an id field and a Map<string, PromiseWithResolvers>.

Servers

Runtime / libraryServer APINotes
BunBun.serve({ websocket })built in, uWebSockets, pub/sub topics
Node + wsnew WebSocketServer()de facto standard, event-emitter API
DenoDeno.upgradeWebSocket(req)returns { socket, response }
Cloudflare Workersnew WebSocketPair()Durable Objects for shared state
Socket.IOnew Server()own protocol on top; not plain WebSocket

Bun

server.upgrade() attaches per-connection data; declare its type with data: {} as T in the websocket handlers (Bun 1.3+).

server.ts
type WsData = { userId: string; rooms: Set<string> };
 
const server = Bun.serve({
  port: 3000,
  fetch(req, server) {
    const url = new URL(req.url);
    if (url.pathname !== "/ws") {
      return new Response("Not found", { status: 404 });
    }
    const userId = userFromCookie(req);
    if (!userId) {
      return new Response("Unauthorized", { status: 401 });
    }
    const ok = server.upgrade(req, {
      data: { userId, rooms: new Set() },
    });
    if (ok) return; // Bun sends the 101 itself
    return new Response("Upgrade required", { status: 426 });
  },
  websocket: {
    data: {} as WsData,
    maxPayloadLength: 64 * 1024, // bytes, default 16 MB
    idleTimeout: 60, // seconds, default 120
    open(ws) {
      ws.subscribe(`user:${ws.data.userId}`);
    },
    message(ws, message) {
      onMessage(ws, message); // string | Buffer
    },
    close(ws, code, reason) {
      console.info(ws.data.userId, "left", code, reason);
    },
  },
});
 
console.info(`listening on ${server.url}`);
ServerWebSocketNotes
ws.datatyped per-connection state from upgrade()
ws.send(data, compress?)returns bytes sent, -1 backpressure, 0 dropped
ws.subscribe(topic)join a pub/sub topic; unsubscribe, isSubscribed
ws.publish(topic, data)send to every subscriber except ws
server.publish(t, data)send to every subscriber
ws.close(code?, reason?)graceful; ws.terminate() drops the socket
ws.remoteAddressclient IP (behind a proxy, read the forwarded header)

Node with ws

node-server.ts
import { createServer } from "node:http";
import { WebSocketServer } from "ws";
 
const http = createServer((_req, res) => {
  res.writeHead(404).end();
});
const wss = new WebSocketServer({
  noServer: true,
  maxPayload: 64 * 1024,
});
 
http.on("upgrade", async (req, socket, head) => {
  const userId = await verifyUpgrade(req); // origin + cookie
  if (!userId) {
    socket.end("HTTP/1.1 401 Unauthorized\r\n\r\n");
    return;
  }
  wss.handleUpgrade(req, socket, head, (ws) => {
    ws.on("error", (err) => console.error(err));
    ws.on("message", (data, isBinary) => {
      if (isBinary) return ws.close(1003, "text only");
      const conn = {
        send: (s: string) => ws.send(s),
        userId,
      };
      onMessage(conn, String(data));
    });
  });
});
 
http.listen(3000);
  • noServer: true + handleUpgrade lets you authenticate before accepting; verifyClient is discouraged by the ws docs.
  • message data is Buffer | ArrayBuffer | Buffer[] (RawData); isBinary tells text from binary.
  • Always attach an error listener: an unhandled error event crashes the process.

Node client

Node 22+ ships the browser-compatible WebSocket as a global (from undici), so the Browser API code above runs unchanged in Node, Bun, Deno and browsers. Use ws only when you need client-side extras such as custom headers or ping().

Heartbeats & reconnection

TCP does not notice a silently dead peer (sleeping laptop, NAT timeout) for minutes. Ping on a timer, drop connections that stop answering, and keep the interval below every proxy's idle timeout (often 60 s).

SideMechanism
Bun serversendPings: true (default) + idleTimeout closes silent sockets
ws serverws.ping() on an interval, terminate() if no pong arrived
Browsercannot send ping frames: send an app-level { type: "ping" }
Load balancerraise idle timeout or ping more often than it
ws-heartbeat.ts
import { WebSocketServer, type WebSocket } from "ws";
 
const wss = new WebSocketServer({ port: 8080 });
const alive = new WeakSet<WebSocket>();
 
wss.on("connection", (ws) => {
  alive.add(ws);
  ws.on("pong", () => alive.add(ws));
});
 
const timer = setInterval(() => {
  for (const ws of wss.clients) {
    if (!alive.has(ws)) {
      ws.terminate();
      continue;
    }
    alive.delete(ws);
    ws.ping();
  }
}, 30_000);
 
wss.on("close", () => clearInterval(timer));

Reconnect with backoff and jitter

"Full jitter" spreads reconnects out so a server restart does not cause a thundering herd.

reconnect.ts
export function backoff(
  attempt: number,
  baseMs = 500,
  capMs = 30_000,
): number {
  return Math.random() * Math.min(
    capMs,
    baseMs * 2 ** attempt,
  );
}
 
type Options = {
  url: string;
  onOpen: (ws: WebSocket) => void; // resubscribe here
  onMessage: (data: unknown) => void;
};
 
export function connect({
  url,
  onOpen,
  onMessage,
}: Options) {
  let attempt = 0;
  let stopped = false;
  let timer: ReturnType<typeof setTimeout> | undefined;
  let ws: WebSocket;
 
  const open = () => {
    ws = new WebSocket(url);
    ws.onopen = () => {
      attempt = 0;
      onOpen(ws);
    };
    ws.onmessage = (e) => onMessage(e.data);
    ws.onclose = (e) => {
      // 1008: auth/policy failure, retrying will not help
      if (stopped || e.code === 1008) return;
      timer = setTimeout(open, backoff(attempt++));
    };
  };
  open();
 
  return {
    send: (data: string) => ws.send(data),
    close() {
      stopped = true;
      clearTimeout(timer);
      ws.close(1000);
    },
  };
}
  • Resubscribe on every open: the server has forgotten your rooms and topics.
  • Resume: number server messages (seq), remember the last one, and send { type: "resume", since: lastSeq } so the server can replay what was missed from a log or stream.
  • Reconnect immediately on the online event and when a hidden tab becomes visible.
  • Queue outgoing messages while disconnected, or reject them; never call send() on a closed socket.

Backpressure

send() never blocks: data piles up in memory if the network is slower than the producer. Watch the queue and pause the producer.

WhereSignalResume when
Browserws.bufferedAmount growspoll until it drops
Bunws.send() returns -1 (queued) or 0 (dropped)drain(ws) handler fires
ws (Node)ws.bufferedAmount; send(data, cb) callbackcallback runs / amount drops
ws streamscreateWebSocketStream(ws) Duplexstandard stream backpressure
client-backpressure.ts
const HIGH_WATER = 1024 * 1024; // 1 MiB
 
async function upload(
  ws: WebSocket,
  chunks: AsyncIterable<Uint8Array<ArrayBuffer>>,
): Promise<void> {
  for await (const chunk of chunks) {
    while (ws.bufferedAmount > HIGH_WATER) {
      await new Promise((r) => setTimeout(r, 50));
    }
    ws.send(chunk);
  }
}
bun-backpressure.ts
type Data = { queue: string[] };
 
function flush(ws: Bun.ServerWebSocket<Data>): void {
  while (ws.data.queue.length > 0) {
    const status = ws.send(ws.data.queue[0]!);
    if (status === -1) return; // queued; wait for drain
    ws.data.queue.shift(); // sent (>0) or dropped (0)
  }
}
 
Bun.serve({
  fetch: (req, server) =>
    server.upgrade(req, { data: { queue: [] } })
      ? undefined
      : new Response("Expected WebSocket", { status: 426 }),
  websocket: {
    data: {} as Data,
    backpressureLimit: 1024 * 1024,
    closeOnBackpressureLimit: true, // shed slow consumers
    message(ws, msg) {
      ws.data.queue.push(String(msg));
      flush(ws);
    },
    drain: flush,
  },
});

Slow consumers are a policy decision: drop messages, coalesce state updates (send only the latest), or disconnect with 1013.

Authentication

The browser WebSocket constructor cannot set an Authorization header, so pick one of:

ApproachHowTrade-off
Session cookiesent automatically with the upgrade requestsimplest; must check Origin (CSWSH)
Token in first messageaccept, then require { type: "auth" } within ~5 sworks cross-origin; socket briefly anonymous
One-time ticket in queryfetch a single-use ticket, open ?ticket=…fine if short-lived and burned on use
Long-lived token in query?token=…avoid: lands in logs, proxies and history
Token in Sec-WebSocket-Protocolsmuggle as a fake subprotocolhacky; also logged; server must echo it

Verify during the upgrade whenever you can, so unauthenticated clients never get a socket:

upgrade-auth.ts
const ALLOWED = new Set(["https://app.example.com"]);
 
Bun.serve({
  async fetch(req, server) {
    const origin = req.headers.get("origin");
    if (!origin || !ALLOWED.has(origin)) {
      return new Response("Forbidden", { status: 403 });
    }
    const cookies = new Bun.CookieMap(
      req.headers.get("cookie") ?? "",
    );
    const session = await sessions.get(
      cookies.get("__Host-sid") ?? "",
    );
    if (!session) return new Response(null, { status: 401 });
    if (server.upgrade(req, { data: session })) return;
    return new Response(null, { status: 426 });
  },
  websocket: {
    data: {} as Session,
    message(ws, msg) {
      // ws.data.userId is trusted from here on
    },
  },
});
  • Authorize every message too (can this user post to this room?), not just the connection.
  • Sessions expire while sockets stay open: re-check periodically or close with 4001 (app-defined) on logout and let the client re-authenticate.
  • See Authentication for cookies, sessions and JWTs.

Scaling

Each server only knows its own sockets. To reach a user connected elsewhere, fan messages out through a broker.

ConcernApproach
One processBun topics: ws.subscribe("room:1"), server.publish("room:1", data)
Many processesRedis pub/sub, NATS or Kafka: publish to the broker, each node relays
Deliverypub/sub is fire-and-forget; use a log (Redis Streams, DB) for replay
Load balancermust forward Upgrade; raise idle timeouts; HTTP/1.1 to the backend
Sticky sessionsa socket stays on one node anyway; stickiness matters for polling fallbacks and in-memory resume state
Connection capsper-user and per-IP limits; raise ulimit -n; budget memory per socket
Deploysclose with 1012, let clients reconnect with jitter to other nodes
redis-fanout.ts
import { RedisClient } from "bun";
 
const pub = new RedisClient(process.env.REDIS_URL);
const sub = await pub.duplicate(); // subscriber conn
 
type Fanout = { topic: string; payload: string };
 
// every node relays broker messages to local sockets
await sub.subscribe("fanout", (message) => {
  const { topic, payload } = JSON.parse(message) as Fanout;
  server.publish(topic, payload);
});
 
export function broadcast(topic: string, payload: string) {
  const msg: Fanout = { topic, payload };
  return pub.publish("fanout", JSON.stringify(msg));
}
# nginx: proxy WebSockets
location /ws {
  proxy_pass http://app;
  proxy_http_version 1.1;
  proxy_set_header Upgrade $http_upgrade;
  proxy_set_header Connection "upgrade";
  proxy_read_timeout 75s;
}

Comparing transports

TransportDirectionDataReconnectSupportReach for it when
WebSocketfull duplextext + binary messagesmanualeverywherechat, multiplayer, collaborative editing
SSE (EventSource)server to clientUTF-8 text eventsautomatic, Last-Event-IDeverywherefeeds, notifications, LLM token streams
Long pollingemulated duplexany HTTP bodyeach requesteverywherefallback through hostile proxies
WebTransportduplex, many streams + datagramsbytes (HTTP/3, QUIC)manualBaseline 2026 (newly available)games, media, unreliable or unordered data
fetch streamingresponse stream; request stream Chromium-onlybytesmanualresponses everywhereone request with a long streamed answer
  • SSE over HTTP/1.1 is limited to 6 connections per origin per browser; HTTP/2 lifts this.
  • WebSocket and SSE both run over TCP and suffer head-of-line blocking; WebTransport does not.
  • Plain HTTP (with caching) beats all of these for data that changes rarely.
  • More on streams: Streaming, Fetch API.

Security

  • Check Origin on every upgrade. WebSockets are not covered by CORS, so a cookie-authenticated socket without an origin allow-list is open to cross-site WebSocket hijacking (CSWSH).
  • wss:// only. Plain ws:// exposes messages and session cookies; Secure cookies are not sent to it anyway.
  • Limit sizes: maxPayloadLength (Bun) / maxPayload (ws); close with 1009 beyond it. Cap permessage-deflate or leave it off: compressed data can expand hugely in memory.
  • Rate limit per connection and per user; close with 1008 or 1013 on abuse.
  • Validate every message with a schema and authorize every action; never trust userId fields sent by the client, use the server-side ws.data.
  • Escape on render: chat text is untrusted input; never inject it as HTML.
  • Bound resources: connections per user, rooms per connection, queued bytes per socket.
rate-limit.ts
// token bucket: `rate` tokens per second,
// bursts up to `burst`
export function tokenBucket(rate: number, burst: number) {
  let tokens = burst;
  let last = performance.now();
  return function take(): boolean {
    const now = performance.now();
    const refill = ((now - last) / 1000) * rate;
    tokens = Math.min(burst, tokens + refill);
    last = now;
    if (tokens < 1) return false;
    tokens -= 1;
    return true;
  };
}
 
// per connection: const allow = tokenBucket(10, 20);
// if (!allow()) ws.close(1008, "rate limit");

Recipes

Wait for the socket to open

When code wants to await a connection before sending instead of nesting sends in an open handler.

export function opened(ws: WebSocket): Promise<void> {
  if (ws.readyState === ws.OPEN) return Promise.resolve();
  return new Promise((resolve, reject) => {
    const off = new AbortController();
    const opts = { once: true, signal: off.signal };
    ws.addEventListener("open", () => {
      off.abort();
      resolve();
    }, opts);
    ws.addEventListener("close", (e) => {
      off.abort();
      reject(new Error(`closed before open: ${e.code}`));
    }, opts);
  });
}
 
const ws = new WebSocket("wss://api.example.com/ws");
await opened(ws);
ws.send(JSON.stringify({ type: "join", room: "lobby" }));

Request and response over one socket

When the client needs a reply to a specific message: tag it with an id, match replies, time out stragglers.

type Call = PromiseWithResolvers<unknown>;
type Reply = { id?: string; result?: unknown; err?: string };
 
export function rpc(ws: WebSocket, timeoutMs = 10_000) {
  const pending = new Map<string, Call>();
  ws.addEventListener("message", (e) => {
    const msg = JSON.parse(String(e.data)) as Reply;
    const call = msg.id ? pending.get(msg.id) : undefined;
    if (!call || !msg.id) return; // a push, not a reply
    pending.delete(msg.id);
    if (msg.err) call.reject(new Error(msg.err));
    else call.resolve(msg.result);
  });
  return (method: string, params?: unknown) => {
    const id = crypto.randomUUID();
    const call = Promise.withResolvers<unknown>();
    pending.set(id, call);
    ws.send(JSON.stringify({ id, method, params }));
    const timer = setTimeout(() => {
      pending.delete(id);
      call.reject(new Error(`${method} timed out`));
    }, timeoutMs);
    return call.promise.finally(() => clearTimeout(timer));
  };
}

Queue messages while offline

When sends can happen before the socket opens or during a reconnect; call attach from the onOpen of the reconnect helper.

export function outbox(max = 100) {
  let queue: string[] = [];
  let socket: WebSocket | undefined;
  return {
    send(data: string): void {
      if (socket?.readyState === WebSocket.OPEN) {
        socket.send(data);
      } else {
        queue = [...queue, data].slice(-max); // keep newest
      }
    },
    attach(ws: WebSocket): void { // call from onOpen
      socket = ws;
      for (const data of queue) ws.send(data);
      queue = [];
    },
  };
}
 
const box = outbox();
box.send(JSON.stringify({ type: "typing" })); // queued
box.attach(new WebSocket("wss://api.example.com/ws"));

Rooms with Bun topics

When clients join and leave channels on one server and Bun's pub/sub should do the fan-out.

type Data = { userId: string };
type Msg =
  | { type: "join" | "leave"; room: string }
  | { type: "say"; room: string; text: string };
 
Bun.serve({
  fetch(req, server) {
    const userId = crypto.randomUUID(); // real app: session
    if (server.upgrade(req, { data: { userId } })) return;
    return new Response("Upgrade required", { status: 426 });
  },
  websocket: {
    data: {} as Data,
    message(ws, raw) {
      const msg = JSON.parse(String(raw)) as Msg; // Zod!
      const room = `room:${msg.room}`;
      if (msg.type === "join") ws.subscribe(room);
      else if (msg.type === "leave") ws.unsubscribe(room);
      else if (msg.type === "say" && ws.isSubscribed(room)) {
        const out = { from: ws.data.userId, text: msg.text };
        ws.publish(room, JSON.stringify(out)); // not to self
      }
    },
  },
});

Across several servers, add the Redis relay from Scaling.

Graceful shutdown

When deploying: tell clients to reconnect elsewhere before the process exits.

const sockets = new Set<Bun.ServerWebSocket<undefined>>();
 
const server = Bun.serve({
  fetch(req, server) {
    if (server.upgrade(req)) return;
    return new Response("Upgrade required", { status: 426 });
  },
  websocket: {
    open: (ws) => void sockets.add(ws),
    close: (ws) => void sockets.delete(ws),
    message(ws, msg) { ws.send(msg); },
  },
});
 
async function shutdown(): Promise<void> {
  // 1012 "service restart": clients reconnect with jitter
  for (const ws of sockets) ws.close(1012, "restarting");
  await server.stop(); // stop accepting, finish in-flight
  process.exit(0);
}
process.once("SIGTERM", shutdown);
process.once("SIGINT", shutdown);

References