Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions docs/1.guide/2.hooks.md
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,18 @@ const hooks = defineHooks({
// Pair with `peer.bufferedAmount`. Not all adapters emit this.
console.log("[ws] drain", peer);
},

ping(peer, data) {
// An application-level ping control frame arrived from the peer.
// Not all adapters emit this.
console.log("[ws] ping", peer, data);
},

pong(peer, data) {
// An application-level pong control frame arrived from the peer,
// typically in reply to `peer.ping()`. Not all adapters emit this.
console.log("[ws] pong", peer, data);
},
});
```

Expand Down
38 changes: 38 additions & 0 deletions docs/1.guide/3.peer.md
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,40 @@ Abruptly close the connection.

To gracefully close the connection, use `peer.close()`.

### `peer.ping(data?)`

Send an application-level WebSocket ping control frame to the client.

Pair with the [`pong` hook](/guide/hooks) — embedding a correlatable payload (e.g. a timestamp) in `data` — to measure round-trip latency:

```ts
import { defineHooks } from "crossws";

const hooks = defineHooks({
message(peer, message) {
if (message.text() === "measure-latency") {
// Embed the send time in the ping payload; the client echoes it back
// in the pong so it can be told apart from keepalive pongs.
peer.ping(`rtt:${Date.now()}`);
}
},
pong(peer, data) {
const payload = new TextDecoder().decode(data);
if (!payload.startsWith("rtt:")) return; // keepalive or unrelated pong
const rtt = Date.now() - Number(payload.slice(4));
console.log(`[ws] round-trip latency: ${rtt}ms`);
},
});
```

Use the [`ping` hook](/guide/hooks) to observe pings the client sends unprompted (e.g. a graphql-ws style client heartbeat).

> [!NOTE]
> The `pong` hook also fires for the pongs the client sends in reply to the adapter's own internal keepalive pings (see [`idleTimeout`](/adapters#idletimeout)), not only your `peer.ping()` calls. To measure RTT reliably, embed a correlatable payload in `peer.ping(data)` and ignore pongs that don't carry it (as above) rather than relying on a shared timestamp that any pong would consume.

> [!NOTE]
> Not all adapters can send a ping frame, or surface inbound ping/pong frames to the `ping`/`pong` hooks. Refer to the [compatibility table](#compatibility).

## Compatibility

| | [Bun][bun] | [Cloudflare][cfw] | [Cloudflare (durable)][cfd] | [Deno][deno] | [Node (ws)][nodews] | [Node (μWebSockets)][nodeuws] | [SSE][sse] |
Expand All @@ -143,6 +177,8 @@ To gracefully close the connection, use `peer.close()`.
| `terminate()` | ✓ | ✓ [^2] | ✓ | ✓ | ✓ | ✓ | ✓ [^2] |
| `bufferedAmount` | ✓ | ⨉ [^8] | ⨉ [^8] | ✓ | ✓ | ✓ | ⨉ [^8] |
| `drain` hook | ✓ | ⨉ [^9] | ⨉ [^9] | ⨉ [^9] | ✓ | ✓ | ⨉ [^9] |
| `ping()` | ✓ | ⨉ [^10] | ⨉ [^10] | ⨉ [^10] | ✓ | ✓ | ⨉ [^10] |
| `ping` / `pong` hooks | ✓ | ⨉ [^10] | ⨉ [^10] | ⨉ [^10] | ✓ | ✓ | ⨉ [^10] |
| `request` | ✓ | ✓ | ✓ [^30] | ✓ | ✓ [^31] | ✓ [^31] | ✓ |
| `remoteAddress` | ✓ | ⨉ | ⨉ | ✓ | ✓ | ✓ | ⨉ |
| `websocket.url` | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ |
Expand Down Expand Up @@ -179,3 +215,5 @@ To gracefully close the connection, use `peer.close()`.
[^8]: The runtime exposes no send-buffer signal, so `peer.bufferedAmount` reports `0`. (Cloudflare buffers and applies backpressure internally; SSE has no equivalent.)

[^9]: The runtime emits no drain signal. Poll [`peer.bufferedAmount`](#peerbufferedamount) instead where it is available.

[^10]: The runtime does not expose a ping/pong control-frame API to user code (pings/pongs are still handled transparently at the protocol level for liveness).
38 changes: 33 additions & 5 deletions src/adapters/bun.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { toBufferLike } from "../utils.ts";
import { adapterUtils, getPeers, DEFAULT_IDLE_TIMEOUT } from "../adapter.ts";
import { AdapterHookable } from "../hooks.ts";
import { Message } from "../message.ts";
import { WSError } from "../error.ts";
import { Peer, type PeerContext } from "../peer.ts";
import type { SyncDriver } from "../sync.ts";

Expand Down Expand Up @@ -65,26 +66,38 @@ const bunAdapter: Adapter<BunAdapter, BunOptions> = (options = {}) => {
idleTimeout: options.idleTimeout ?? DEFAULT_IDLE_TIMEOUT,
message: (ws, message) => {
const peers = getPeers(globalPeers, ws.data.namespace);
const peer = getPeer(ws, peers, baseUtils.sync);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("message", peer, new Message(message, peer));
},
open: (ws) => {
const peers = getPeers(globalPeers, ws.data.namespace);
const peer = getPeer(ws, peers, baseUtils.sync);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
peers.add(peer);
hooks.callHook("open", peer);
},
close: (ws, code, reason) => {
const peers = getPeers(globalPeers, ws.data.namespace);
const peer = getPeer(ws, peers, baseUtils.sync);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
peers.delete(peer);
hooks.callHook("close", peer, { code, reason });
},
drain: (ws) => {
const peers = getPeers(globalPeers, ws.data.namespace);
const peer = getPeer(ws, peers);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("drain", peer);
},
// Bun auto-replies to an inbound ping with a pong per the spec; these
// hooks only observe the control frames, they don't need to answer them.
ping: (ws, data) => {
const peers = getPeers(globalPeers, ws.data.namespace);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("ping", peer, data);
},
pong: (ws, data) => {
const peers = getPeers(globalPeers, ws.data.namespace);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("pong", peer, data);
},
},
};
};
Expand All @@ -96,7 +109,8 @@ export default bunAdapter;
function getPeer(
ws: ServerWebSocket<ContextData>,
peers: Set<BunPeer>,
sync?: SyncDriver,
sync: SyncDriver | undefined,
hooks: AdapterHookable,
): BunPeer {
if (ws.data.peer) {
return ws.data.peer;
Expand All @@ -107,6 +121,7 @@ function getPeer(
peers,
namespace: ws.data.namespace,
sync,
hooks,
});
ws.data.peer = peer;
return peer;
Expand All @@ -118,6 +133,7 @@ class BunPeer extends Peer<{
request: Request;
peers: Set<BunPeer>;
sync?: SyncDriver;
hooks: AdapterHookable;
}> {
override get remoteAddress(): string {
return this._internal.ws.remoteAddress;
Expand Down Expand Up @@ -156,4 +172,16 @@ class BunPeer extends Peer<{
override terminate(): void {
this._internal.ws.terminate();
}

override ping(data?: unknown): number {
// Guard against the native ping rejecting the payload (e.g. the 125-byte
// control-frame limit): surface it through the `error` hook rather than
// letting it crash a caller inside a hook handler.
try {
return this._internal.ws.ping(data as any);
} catch (error) {
this._internal.hooks.callHook("error", this, new WSError(error));
return 0;
}
}
}
36 changes: 35 additions & 1 deletion src/adapters/node.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,12 @@ type AugmentedReq = IncomingMessage & {
// `ws` instance tagged with the heartbeat liveness flag (see the idle sweep below).
type HeartbeatWS = WebSocketT & { _isAlive?: boolean };

// Payload carried by our own liveness probe (see the idle sweep). The client
// echoes it back in the pong, letting the `pong` handler tell an internal
// keepalive reply apart from an app-level `peer.ping()` and skip the `pong`
// hook for it. Kept well under the 125-byte control-frame limit.
const HEARTBEAT_PING = Buffer.from("crossws-ping");

export interface NodeAdapter extends AdapterInstance {
handleUpgrade(
req: IncomingMessage,
Expand Down Expand Up @@ -89,6 +95,7 @@ const nodeAdapter: Adapter<NodeAdapter, NodeOptions> = (options = {}) => {
nodeReq,
namespace: nodeReq._namespace,
sync: baseUtils.sync,
hooks,
});
peers.add(peer);
liveSockets.add(ws as HeartbeatWS);
Expand Down Expand Up @@ -116,6 +123,20 @@ const nodeAdapter: Adapter<NodeAdapter, NodeOptions> = (options = {}) => {
}
hooks.callHook("message", peer, new Message(data, peer));
});
// `ws` auto-replies to an inbound ping with a pong per the spec; these
// hooks only observe the control frames, they don't need to answer them.
ws.on("ping", (data: Buffer) => {
hooks.callHook("ping", peer, data);
});
ws.on("pong", (data: Buffer) => {
// Skip our own liveness probe's echoed pong (see the idle sweep); it is
// not an app-level pong and would otherwise fire the `pong` hook every
// idle interval with a bogus payload.
if (data.equals(HEARTBEAT_PING)) {
return;
}
hooks.callHook("pong", peer, data);
});
ws.on("error", (error: Error) => {
peers.delete(peer);
hooks.callHook("error", peer, new WSError(error));
Expand Down Expand Up @@ -160,7 +181,7 @@ const nodeAdapter: Adapter<NodeAdapter, NodeOptions> = (options = {}) => {
}
ws._isAlive = false;
try {
ws.ping();
ws.ping(HEARTBEAT_PING);
} catch {
// socket may have raced into CLOSING between the sweep and the ping
}
Expand Down Expand Up @@ -245,6 +266,7 @@ class NodePeer extends Peer<{
nodeReq: IncomingMessage;
ws: WebSocketT & { _peer?: NodePeer };
sync?: SyncDriver;
hooks: AdapterHookable;
}> {
override get remoteAddress() {
return this._internal.nodeReq.socket?.remoteAddress;
Expand Down Expand Up @@ -290,6 +312,18 @@ class NodePeer extends Peer<{
override terminate() {
this._internal.ws.terminate();
}

override ping(data?: unknown): void {
// `ws` validates the control frame synchronously (rejecting non-buffer
// payloads and anything over the 125-byte limit) and throws. Surface that
// through the `error` hook instead of letting it crash the process — e.g.
// when `ping()` is called from inside a hook handler.
try {
this._internal.ws.ping(data);
} catch (error) {
this._internal.hooks.callHook("error", this, new WSError(error));
}
}
}

// --- web compat ---
Expand Down
49 changes: 44 additions & 5 deletions src/adapters/uws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { toBufferLike } from "../utils.ts";
import { adapterUtils, getPeers, DEFAULT_IDLE_TIMEOUT } from "../adapter.ts";
import { AdapterHookable } from "../hooks.ts";
import { Message } from "../message.ts";
import { WSError } from "../error.ts";
import { Peer, type PeerContext } from "../peer.ts";
import type { SyncDriver } from "../sync.ts";
import { StubRequest } from "../_request.ts";
Expand Down Expand Up @@ -54,7 +55,7 @@ const uwsAdapter: Adapter<UWSAdapter, UWSOptions> = (options = {}) => {
...options.uws,
close(ws, code, message) {
const peers = getPeers(globalPeers, ws.getUserData().namespace);
const peer = getPeer(ws, peers, baseUtils.sync);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
((peer as any)._internal.ws as UwsWebSocketProxy).readyState = 2 /* CLOSING */;
peers.delete(peer);
hooks.callHook("close", peer, {
Expand All @@ -65,17 +66,36 @@ const uwsAdapter: Adapter<UWSAdapter, UWSOptions> = (options = {}) => {
},
message(ws, message, _isBinary) {
const peers = getPeers(globalPeers, ws.getUserData().namespace);
const peer = getPeer(ws, peers, baseUtils.sync);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("message", peer, new Message(message, peer));
},
drain(ws) {
const peers = getPeers(globalPeers, ws.getUserData().namespace);
const peer = getPeer(ws, peers);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("drain", peer);
},
// uWS auto-replies to an inbound ping with a pong per the spec; these
// hooks only observe the control frames, they don't need to answer them.
// Skip the `Uint8Array` copy entirely when nothing consumes it.
ping(ws, message) {
if (!hooks.options.hooks?.ping && !hooks.options.resolve) {
return;
}
const peers = getPeers(globalPeers, ws.getUserData().namespace);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("ping", peer, new Uint8Array(message));
},
pong(ws, message) {
if (!hooks.options.hooks?.pong && !hooks.options.resolve) {
return;
}
const peers = getPeers(globalPeers, ws.getUserData().namespace);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
hooks.callHook("pong", peer, new Uint8Array(message));
},
open(ws) {
const peers = getPeers(globalPeers, ws.getUserData().namespace);
const peer = getPeer(ws, peers, baseUtils.sync);
const peer = getPeer(ws, peers, baseUtils.sync, hooks);
peers.add(peer);
hooks.callHook("open", peer);
},
Expand Down Expand Up @@ -147,7 +167,12 @@ export default uwsAdapter;

// --- peer ---

function getPeer(uws: uws.WebSocket<UserData>, peers: Set<UWSPeer>, sync?: SyncDriver): UWSPeer {
function getPeer(
uws: uws.WebSocket<UserData>,
peers: Set<UWSPeer>,
sync: SyncDriver | undefined,
hooks: AdapterHookable,
): UWSPeer {
const uwsData = uws.getUserData();
if (uwsData.peer) {
return uwsData.peer;
Expand All @@ -160,6 +185,7 @@ function getPeer(uws: uws.WebSocket<UserData>, peers: Set<UWSPeer>, sync?: SyncD
namespace: uwsData.namespace,
uwsData,
sync,
hooks,
});
uwsData.peer = peer;
return peer;
Expand All @@ -173,6 +199,7 @@ class UWSPeer extends Peer<{
ws: UwsWebSocketProxy;
uwsData: UserData;
sync?: SyncDriver;
hooks: AdapterHookable;
}> {
override get remoteAddress(): string | undefined {
try {
Expand Down Expand Up @@ -217,6 +244,18 @@ class UWSPeer extends Peer<{
override terminate(): void {
this._internal.uws.close();
}

override ping(data?: uws.RecognizedString): number {
// Guard against uWS rejecting the payload (e.g. the 125-byte control-frame
// limit): surface it through the `error` hook rather than letting it crash
// a caller inside a hook handler.
try {
return this._internal.uws.ping(data);
} catch (error) {
this._internal.hooks.callHook("error", this, new WSError(error));
return 0;
}
}
}

// --- web compat ---
Expand Down
20 changes: 20 additions & 0 deletions src/hooks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -214,4 +214,24 @@ export interface Hooks {

/** An error occurs */
error: (peer: Peer, error: WSError) => MaybePromise<void>;

/**
* An application-level WebSocket ping control frame was received from the
* peer (e.g. sent by the client, or by another server via
* {@link Peer.ping}).
*
* **Note:** Only emitted by adapters that surface inbound ping frames.
* Refer to the [compatibility table](https://crossws.h3.dev/guide/peer#compatibility).
*/
ping: (peer: Peer, data: Uint8Array) => MaybePromise<void>;

/**
* An application-level WebSocket pong control frame was received from the
* peer, typically in reply to {@link Peer.ping}. Use together with a
* timestamp embedded in the ping payload to measure round-trip latency.
*
* **Note:** Only emitted by adapters that surface inbound pong frames.
* Refer to the [compatibility table](https://crossws.h3.dev/guide/peer#compatibility).
*/
pong: (peer: Peer, data: Uint8Array) => MaybePromise<void>;
}
Loading
Loading