Skip to content

Latest commit

 

History

15 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

WebSocket

A zero-dependency, fully-typed implementation of the WebSocket protocol (RFC 6455) in TypeScript. Built entirely on Node.js built-ins — no runtime dependencies.

Supports per-message deflate compression (RFC 7692), SOCKS5 / HTTP CONNECT proxying, prepared (cached) messages for broadcasting, and message joining.

Installation

npm install
npm run build

Quick Start

Echo Server + Client

import http from 'node:http';
import { Upgrader, Dialer, TextMessage } from './dist/index.js';

// Server
const upgrader = new Upgrader({ enableCompression: true });
const server = http.createServer();

server.on('upgrade', (req, socket, head) => {
  const conn = upgrader.upgrade(req, socket, head);
  void echoLoop(conn);
});

server.listen(8080);

async function echoLoop(conn) {
  for (;;) {
    const { messageType, data } = await conn.readMessage();
    await conn.writeMessage(messageType, data);
  }
}

// Client
const dialer = new Dialer({ handshakeTimeout: 5000 });
const { conn } = await dialer.dial('ws://localhost:8080/');
await conn.writeMessage(TextMessage, 'Hello!');
const result = await conn.readMessage();
console.log(result.data.toString()); // "Hello!"

Core Concepts

Concurrency Model

A Conn supports one concurrent reader and one concurrent writer. This maps naturally to Node's single-threaded event loop — in practice, you run one async function for reads and one for writes.

close() and writeControl() are safe to call concurrently with everything else.

Message Types

Constant Value Description
TextMessage 1 UTF-8 text payload
BinaryMessage 2 Binary payload
CloseMessage 8 Connection close
PingMessage 9 Ping (keepalive)
PongMessage 10 Pong response

These are exported from both the top-level index and individually from each source module.


Server API

Upgrader

Performs the HTTP → WebSocket upgrade handshake. Use it inside Node's server.on('upgrade') handler.

import http from 'node:http';
import { Upgrader } from 'websocket';

const upgrader = new Upgrader({
  readBufferSize: 4096,
  writeBufferSize: 4096,
  subprotocols: ['chat', 'superchat'],
  enableCompression: true,
  checkOrigin(req) {
    return req.headers.origin === 'https://myapp.com';
  },
});

const server = http.createServer();
server.on('upgrade', (req, socket, head) => {
  const conn = upgrader.upgrade(req, socket, head, {
    'Set-Cookie': 'session=abc123',
  });
  // conn.subprotocol — the negotiated subprotocol, if any
  // conn.isServer === true
});

UpgraderOptions

Option Type Default Description
handshakeTimeout number 0 (none) Timeout in ms for the handshake write
readBufferSize number 4096 Internal read buffer size
writeBufferSize number 4096 Internal write buffer size (header space is added automatically)
subprotocols string[] [] Supported subprotocols in preference order
checkOrigin (req) => boolean same-origin Origin validation function; default rejects cross-origin
enableCompression boolean false Negotiate permessage-deflate
errorHandler (res, status, reason) => void Custom HTTP error response for failed handshakes

upgrade(req, socket, head, responseHeaders?)

Param Type Description
req http.IncomingMessage The upgrade request
socket net.Socket The raw TCP socket
head Buffer First packet of the new stream (may be empty)
responseHeaders Record<string, string> Additional headers in the 101 response (e.g. cookies)

Throws HandshakeError if the request is not a valid WebSocket upgrade.

isWebSocketUpgrade(req)

Convenience predicate: returns true if the request has Connection: Upgrade and Upgrade: websocket.

import { isWebSocketUpgrade, subprotocols } from 'websocket';

if (isWebSocketUpgrade(req)) {
  const requested = subprotocols(req); // ['chat', 'superchat']
}

Client API

Dialer

Connects to a WebSocket server by performing the client-side upgrade handshake.

import { Dialer, TextMessage } from 'websocket';

const dialer = new Dialer({
  handshakeTimeout: 10000,
  subprotocols: ['chat'],
  enableCompression: true,
  headers: { 'Authorization': 'Bearer token' },
});

const { conn, resp } = await dialer.dial('wss://echo.example.com/');
// resp.statusCode === 101
// resp.headers — response headers from the upgrade
// conn.subprotocol — the negotiated subprotocol

DialerOptions

Option Type Default Description
handshakeTimeout number 45000 Timeout in ms
readBufferSize number 4096 Internal read buffer size
writeBufferSize number 4096 Internal write buffer size
subprotocols string[] [] Client subprotocol preferences
enableCompression boolean false Advertise permessage-deflate support
tlsConfig tls.ConnectionOptions TLS settings for wss:// connections
headers Record<string, string> {} Extra HTTP headers on the upgrade request
proxy (req) => Promise<string | undefined> Proxy resolution function

DefaultDialer

A pre-configured singleton using Proxy: http.ProxyFromEnvironment, 45s timeout, and default buffer sizes.

import { DefaultDialer } from 'websocket';
const { conn } = await DefaultDialer.dial('ws://localhost:8080/');

Connection API

Reading Messages

// Read a complete message into a Buffer
const { messageType, data } = await conn.readMessage();
console.log(messageType); // TextMessage (1) or BinaryMessage (2)

// Streaming read (for large messages)
const { messageType, reader } = await conn.nextReader();
const chunks = [];
for (;;) {
  const chunk = await reader.read();
  if (chunk === null) break;
  chunks.push(chunk);
}

Writing Messages

// Write a complete message
await conn.writeMessage(TextMessage, 'Hello');
await conn.writeMessage(BinaryMessage, Buffer.from([0x00, 0x01]));

// Stream a message (fragments automatically)
const writer = conn.nextWriter(TextMessage);
writer.write('part 1');
writer.write('part 2');
await writer.close();

Control Frames

// Send a ping (keepalive)
conn.writeControl(PingMessage, Buffer.from('heartbeat'));

// Send a graceful close
conn.writeControl(CloseMessage, FormatCloseMessage(1000, 'bye'));

// Force-close the underlying socket (no close handshake)
conn.close();

Connection State

conn.isServer            // true if server-side

// Deadlines (applied per-read / per-frame-write)
conn.setReadDeadline(30000);   // 30s read timeout
conn.setWriteDeadline(10000);  // 10s write timeout

// Read limit (auto-closes with 1009 if exceeded)
conn.setReadLimit(1024 * 1024); // 1 MB max message

// Network info
conn.localAddr   // { address, family, port }
conn.remoteAddr  // { address, family, port }

Control Handlers

Customize how the connection responds to control frames:

// Close handler (default: echoes the close code back)
conn.setCloseHandler((code, text) => {
  console.log(`Peer closing: ${code} ${text}`);
});

// Ping handler (default: auto-replies with pong)
conn.setPingHandler((data) => {
  console.log('Ping received:', data);
});

// Pong handler (default: no-op)
conn.setPongHandler((data) => {
  // Measure latency using custom data payloads
  const elapsed = Date.now() - parseInt(data, 10);
  console.log(`Latency: ${elapsed}ms`);
});

Compression

Per-message deflate (RFC 7692) in "no context takeover" mode (fresh compression context per message).

// Server side — enable in Upgrader options
const upgrader = new Upgrader({ enableCompression: true });

// Client side — enable in Dialer options
const dialer = new Dialer({ enableCompression: true });

// Per-connection toggling
conn.enableWriteCompression = false;
conn.setCompressionLevel(6); // -2 (HuffmanOnly) to 9 (BestCompression)

Compression is negotiated during the handshake. Both sides must advertise support. Once negotiated, read decompression is automatic. Write compression can be toggled per-message via conn.enableWriteCompression.

If you need to integrate compression into custom flows:

import { compressNoContextTakeover, decompressNoContextTakeover } from 'websocket';

// Wrap a writer
const compressedWriter = compressNoContextTakeover(rawWriter, 1);

// Wrap a reader
const decompressedReader = decompressNoContextTakeover(rawReader);

Proxy Support

SOCKS5 and HTTP CONNECT proxies are supported on the client dial path.

SOCKS5

import { createProxyDialer } from 'websocket';

const forwardDial = (host, port) => {
  // standard TCP dial
};
const proxyDial = createProxyDialer('socks5://user:pass@proxy:1080', forwardDial);
const sock = await proxyDial('echo.example.com', 80);

HTTP CONNECT

const proxyDial = createProxyDialer('http://proxy.corp:8080', forwardDial);
const sock = await proxyDial('echo.example.com', 443);

Environment-Variable Proxying

Use the DefaultDialer or implement DialerOptions.proxy to resolve proxy URLs dynamically (e.g. from ALL_PROXY, NO_PROXY environment variables).


JSON Helpers

Convenience wrappers for JSON serialization over WebSocket text messages:

import { writeJSON, readJSON } from 'websocket';

await writeJSON(conn, { type: 'greeting', body: 'hello' });
const msg = await readJSON(conn);
// msg === { type: 'greeting', body: 'hello' }

Prepared Messages

Cache wire-format frame data for broadcasting the same message to many connections. Each unique combination of (server/client, compress, level) is encoded once:

import { PreparedMessage, TextMessage } from 'websocket';

const pm = new PreparedMessage(TextMessage, JSON.stringify({ event: 'tick', ts: Date.now() }));

// Broadcast to many connections without re-encoding
for (const conn of connections) {
  await conn.writePreparedMessage(pm);
}

Message Joining

Concatenate consecutive WebSocket messages into a single readable stream with a delimiter:

import { joinMessages } from 'websocket';

const reader = joinMessages(conn, '\n');
const chunks = [];
for (;;) {
  const chunk = await reader.read();
  if (chunk === null) break;
  chunks.push(chunk);
}
const full = Buffer.concat(chunks).toString();

Error Handling

CloseError

Thrown when the peer sends a close frame. Carries the close code and optional text:

import { CloseError, isCloseError, isUnexpectedCloseError } from 'websocket';

try {
  await conn.readMessage();
} catch (err) {
  if (isCloseError(err)) {
    console.log(`Closed: ${err.code}`);
  }
  if (isUnexpectedCloseError(err, 1000, 1001)) {
    // Not a normal closure — log and investigate
    console.error('Abnormal close:', err);
  }
}

Close Status Codes

Constant Code Meaning
CloseNormalClosure 1000 Normal closure
CloseGoingAway 1001 Endpoint going away
CloseProtocolError 1002 Protocol error
CloseUnsupportedData 1003 Received data type not supported
CloseNoStatusReceived 1005 No status (reserved)
CloseAbnormalClosure 1006 Abnormal (reserved)
CloseInvalidFramePayloadData 1007 Invalid payload data
ClosePolicyViolation 1008 Policy violation
CloseMessageTooBig 1009 Message too big
CloseMandatoryExtension 1010 Extension expected
CloseInternalServerErr 1011 Internal server error
CloseServiceRestart 1012 Service restart
CloseTryAgainLater 1013 Try again later
CloseTLSHandshake 1015 TLS handshake failure

Sentinel Errors

Error When
ErrCloseSent Write attempted after a close frame was sent
ErrReadLimit Message exceeds setReadLimit

FormatCloseMessage(code, text?)

Builds a close frame payload (2-byte big-endian code + UTF-8 text):

import { FormatCloseMessage, CloseNormalClosure } from 'websocket';
conn.writeControl(CloseMessage, FormatCloseMessage(CloseNormalClosure, 'bye'));

Protocol-Level Utilities

For advanced use cases (custom frame handling, protocol debugging):

import {
  computeAcceptKey,       // SHA-1(challengeKey + GUID), base64
  generateChallengeKey,   // 16 random bytes, base64
  parseExtensions,        // Parse Sec-WebSocket-Extensions header
  parseFrameHeader,       // Parse raw frame header bytes into structured data
  maskBytes,              // XOR bytes with rotating 4-byte key
  tokenListContainsValue, // RFC 2616 1#token header value check
} from 'websocket';

const key = computeAcceptKey('dGhlIHNhbXBsZSBub25jZQ==');
// "s3pPLMBiTxaQ9kYGzzhZRbK+xOo="

Complete Example

import http from 'node:http';
import {
  Upgrader, Dialer,
  Conn,
  TextMessage, CloseMessage,
  CloseNormalClosure, CloseError,
  FormatCloseMessage, isCloseError,
} from 'websocket';

// ── Server ────────────────────────────────────────────
const upgrader = new Upgrader({ enableCompression: true });

const server = http.createServer();
server.on('upgrade', (req, socket, head) => {
  const conn = upgrader.upgrade(req, socket, head);
  console.log('New connection, subprotocol:', conn.subprotocol);
  handleClient(conn).catch(() => {});
});

async function handleClient(conn: Conn) {
  try {
    for (;;) {
      const { messageType, data } = await conn.readMessage();
      console.log('Received:', data.toString());
      await conn.writeMessage(messageType, data); // echo
    }
  } catch (err) {
    if (isCloseError(err)) {
      console.log('Client closed:', err.code);
    }
  }
}

server.listen(8080, () => console.log('Listening on :8080'));

// ── Client ────────────────────────────────────────────
const dialer = new Dialer({
  handshakeTimeout: 5000,
  enableCompression: true,
});

const { conn } = await dialer.dial('ws://localhost:8080/');

// Send a message
await conn.writeMessage(TextMessage, 'Hello from client!');

// Read the echo
const { data } = await conn.readMessage();
console.log('Echo:', data.toString());

// Graceful close
await conn.writeControl(CloseMessage, FormatCloseMessage(CloseNormalClosure, 'done'));

Testing

npm test          # all tests (75)
# or individually:
npm run test:e2e  # end-to-end integration only

Built with Node's native test runner. No test framework dependencies. Tests cover frame parsing, masking, compression round-trips, server upgrade validation, client handshake, proxying, JSON serialization, prepared message caching, message joining, and end-to-end echo flows.

About

Implementation of the WebSocket Protocol

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages