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| Header | Sent by | Purpose |
|---|---|---|
Upgrade: websocket | both | request / confirm the protocol switch |
Connection: Upgrade | both | marks Upgrade as hop-by-hop |
Sec-WebSocket-Version | client | always 13; server answers 426 with supported versions |
Sec-WebSocket-Key | client | 16 random bytes, base64; proves the server speaks WebSocket |
Sec-WebSocket-Accept | server | base64 of SHA-1 of key + fixed GUID |
Origin | browser | page origin; the server must check it (no CORS here) |
Sec-WebSocket-Protocol | both | client offers subprotocols, server picks at most one |
Sec-WebSocket-Extensions | both | negotiated extensions, e.g. permessage-deflate |
The accept value is not security, only proof that the server understood the upgrade:
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
| Field | Bits | Meaning |
|---|---|---|
| FIN | 1 | last frame of a message |
| RSV1–3 | 3 | extension flags (RSV1 = compressed by permessage-deflate) |
| opcode | 4 | frame type, see below |
| MASK | 1 | payload is masked; must be 1 client to server |
| Payload len | 7, +16, +64 | 0–125 literal; 126: next 16 bits; 127: next 64 bits |
| Masking key | 0 or 32 | present when MASK = 1 |
| Payload | n | extension data + application data |
| Opcode | Frame | Notes |
|---|---|---|
0x0 | continuation | later fragments of a fragmented message |
0x1 | text | payload must be valid UTF-8 (else close 1007) |
0x2 | binary | arbitrary bytes |
0x3–0x7 | reserved | future data frames |
0x8 | close | optional 2-byte code + UTF-8 reason |
0x9 | ping | peer must answer with a pong carrying the same payload |
0xA | pong | reply to ping, or unsolicited one-way heartbeat |
0xB–0xF | reserved | future 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
| Scheme | Port | Transport | Use |
|---|---|---|---|
ws:// | 80 | plain TCP | local development only |
wss:// | 443 | TLS over TCP | always in production; proxies pass it through unmodified |
- An
https:page cannot openws://(mixed content). Relativehttp(s):URLs are accepted by the constructor and mapped tows(s):in current browsers. - Subprotocols are application-level version tags:
new WebSocket(url, ["chat.v2", "chat.v1"]). The server echoes one;ws.protocolis""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
| Code | Name | Meaning |
|---|---|---|
1000 | Normal closure | done, purpose fulfilled |
1001 | Going away | page navigated away, server shutting down |
1002 | Protocol error | malformed frame, unmasked client frame |
1003 | Unsupported data | e.g. binary sent to a text-only endpoint |
1005 | No status received | reserved: close frame had no code (never sent) |
1006 | Abnormal closure | reserved: no close frame, e.g. TCP dropped |
1007 | Invalid payload | text frame was not UTF-8 |
1008 | Policy violation | generic "not allowed" (auth failed, bad origin) |
1009 | Message too big | exceeded the peer's size limit |
1010 | Mandatory extension | client needed an extension the server refused |
1011 | Internal error | unexpected server condition |
1012 | Service restart | reconnect later (IANA registry) |
1013 | Try again later | overloaded or rate limited (IANA registry) |
1015 | TLS handshake failure | reserved: reported locally only |
3000–3999 | registered | libraries and frameworks, via IANA |
4000–4999 | private | your application's own codes |
Browser API
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
| Value | Constant | Meaning | send() |
|---|---|---|---|
0 | WebSocket.CONNECTING | handshake in progress | throws InvalidStateError |
1 | WebSocket.OPEN | ready | queues the data |
2 | WebSocket.CLOSING | close handshake started | silently discarded |
3 | WebSocket.CLOSED | closed or failed to open | silently discarded |
Members
| Member | Type | Notes |
|---|---|---|
new WebSocket(url, protocols?) | WebSocket | connects immediately; no custom headers possible |
send(data) | void | string, ArrayBuffer, Blob, TypedArray, DataView |
close(code?, reason?) | void | code 1000 or 3000–4999; reason max 123 UTF-8 bytes |
binaryType | "blob" | "arraybuffer" | type of binary e.data |
bufferedAmount | number | bytes queued by send() but not yet sent |
readyState | 0 | 1 | 2 | 3 | see above |
protocol, extensions | string | what the server accepted |
url | string | resolved absolute URL |
onopen … onclose | handler props | or addEventListener |
| Event | Object | Useful fields |
|---|---|---|
open | Event | none |
message | MessageEvent | data: string | Blob | ArrayBuffer, origin |
error | Event | none; deliberately opaque |
close | CloseEvent | code, reason, wasClean |
- Calling
close()with any other code throwsInvalidAccessError; a longer reason throwsSyntaxError. - 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.
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.
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.tsbetween 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
idfield and aMap<string, PromiseWithResolvers>.
Servers
| Runtime / library | Server API | Notes |
|---|---|---|
| Bun | Bun.serve({ websocket }) | built in, uWebSockets, pub/sub topics |
Node + ws | new WebSocketServer() | de facto standard, event-emitter API |
| Deno | Deno.upgradeWebSocket(req) | returns { socket, response } |
| Cloudflare Workers | new WebSocketPair() | Durable Objects for shared state |
| Socket.IO | new 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+).
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}`);ServerWebSocket | Notes |
|---|---|
ws.data | typed 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.remoteAddress | client IP (behind a proxy, read the forwarded header) |
Node with ws
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+handleUpgradelets you authenticate before accepting;verifyClientis discouraged by thewsdocs.messagedata isBuffer | ArrayBuffer | Buffer[](RawData);isBinarytells text from binary.- Always attach an
errorlistener: an unhandlederrorevent 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).
| Side | Mechanism |
|---|---|
| Bun server | sendPings: true (default) + idleTimeout closes silent sockets |
ws server | ws.ping() on an interval, terminate() if no pong arrived |
| Browser | cannot send ping frames: send an app-level { type: "ping" } |
| Load balancer | raise idle timeout or ping more often than it |
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.
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
onlineevent 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.
| Where | Signal | Resume when |
|---|---|---|
| Browser | ws.bufferedAmount grows | poll until it drops |
| Bun | ws.send() returns -1 (queued) or 0 (dropped) | drain(ws) handler fires |
ws (Node) | ws.bufferedAmount; send(data, cb) callback | callback runs / amount drops |
ws streams | createWebSocketStream(ws) Duplex | standard stream backpressure |
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);
}
}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:
| Approach | How | Trade-off |
|---|---|---|
| Session cookie | sent automatically with the upgrade request | simplest; must check Origin (CSWSH) |
| Token in first message | accept, then require { type: "auth" } within ~5 s | works cross-origin; socket briefly anonymous |
| One-time ticket in query | fetch 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-Protocol | smuggle as a fake subprotocol | hacky; also logged; server must echo it |
Verify during the upgrade whenever you can, so unauthenticated clients never get a socket:
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.
| Concern | Approach |
|---|---|
| One process | Bun topics: ws.subscribe("room:1"), server.publish("room:1", data) |
| Many processes | Redis pub/sub, NATS or Kafka: publish to the broker, each node relays |
| Delivery | pub/sub is fire-and-forget; use a log (Redis Streams, DB) for replay |
| Load balancer | must forward Upgrade; raise idle timeouts; HTTP/1.1 to the backend |
| Sticky sessions | a socket stays on one node anyway; stickiness matters for polling fallbacks and in-memory resume state |
| Connection caps | per-user and per-IP limits; raise ulimit -n; budget memory per socket |
| Deploys | close with 1012, let clients reconnect with jitter to other nodes |
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
| Transport | Direction | Data | Reconnect | Support | Reach for it when |
|---|---|---|---|---|---|
| WebSocket | full duplex | text + binary messages | manual | everywhere | chat, multiplayer, collaborative editing |
SSE (EventSource) | server to client | UTF-8 text events | automatic, Last-Event-ID | everywhere | feeds, notifications, LLM token streams |
| Long polling | emulated duplex | any HTTP body | each request | everywhere | fallback through hostile proxies |
| WebTransport | duplex, many streams + datagrams | bytes (HTTP/3, QUIC) | manual | Baseline 2026 (newly available) | games, media, unreliable or unordered data |
| fetch streaming | response stream; request stream Chromium-only | bytes | manual | responses everywhere | one 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
Originon 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. Plainws://exposes messages and session cookies;Securecookies are not sent to it anyway.- Limit sizes:
maxPayloadLength(Bun) /maxPayload(ws); close with1009beyond it. Cappermessage-deflateor leave it off: compressed data can expand hugely in memory. - Rate limit per connection and per user; close with
1008or1013on abuse. - Validate every message with a schema and authorize every action; never trust
userIdfields sent by the client, use the server-sidews.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.
// 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
- MDN: WebSocket (opens in a new tab), Writing WebSocket servers (opens in a new tab), WebSocketStream (opens in a new tab), WebTransport (opens in a new tab)
- High Performance Browser Networking: WebSocket (opens in a new tab)
- RFC 6455: The WebSocket Protocol (opens in a new tab), IANA close code registry (opens in a new tab)
- Bun: WebSockets (opens in a new tab)
- ws on GitHub (opens in a new tab)
- OWASP WebSocket Security Cheat Sheet (opens in a new tab)