From e44e11c5201d6d3da06a905d658628bfccf0c0cd Mon Sep 17 00:00:00 2001 From: Pooya Parsa Date: Fri, 3 Jul 2026 09:48:52 +0000 Subject: [PATCH 1/2] feat: application-level ping/pong hooks and peer.ping() Closes the remaining follow-up from #154 after #201 landed automatic liveness. Adds `ping`/`pong` hooks to observe inbound control frames and `peer.ping(data?)` to send one, wired for Node (ws), uWebSockets, and Bun (all natively support it); Deno, Cloudflare, Bunny, and SSE have no such runtime API and fall back to a warning no-op. Co-Authored-By: Claude Sonnet 5 --- docs/1.guide/2.hooks.md | 12 +++++++ docs/1.guide/3.peer.md | 32 ++++++++++++++++++ src/adapters/bun.ts | 16 +++++++++ src/adapters/node.ts | 12 +++++++ src/adapters/uws.ts | 16 +++++++++ src/hooks.ts | 20 +++++++++++ src/peer.ts | 15 +++++++++ src/server/_resolve.ts | 11 +++++- test/adapters/bun.test.ts | 6 +++- test/adapters/node.test.ts | 4 ++- test/adapters/uws.test.ts | 4 ++- test/fixture/_shared.ts | 10 ++++++ test/tests.ts | 69 ++++++++++++++++++++++++++++++++++++++ 13 files changed, 223 insertions(+), 4 deletions(-) diff --git a/docs/1.guide/2.hooks.md b/docs/1.guide/2.hooks.md index d4921c9..c4dbc7b 100644 --- a/docs/1.guide/2.hooks.md +++ b/docs/1.guide/2.hooks.md @@ -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); + }, }); ``` diff --git a/docs/1.guide/3.peer.md b/docs/1.guide/3.peer.md index c1b85c9..959822b 100644 --- a/docs/1.guide/3.peer.md +++ b/docs/1.guide/3.peer.md @@ -133,6 +133,34 @@ 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) — e.g. embedding 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") { + peer.context.pingSentAt = Date.now(); + peer.ping(); + } + }, + pong(peer) { + const rtt = Date.now() - (peer.context.pingSentAt as number); + 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] +> 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] | @@ -143,6 +171,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` | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ | @@ -179,3 +209,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). diff --git a/src/adapters/bun.ts b/src/adapters/bun.ts index 8f2f440..2f86c02 100644 --- a/src/adapters/bun.ts +++ b/src/adapters/bun.ts @@ -85,6 +85,18 @@ const bunAdapter: Adapter = (options = {}) => { const peer = getPeer(ws, peers); 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.callHook("ping", peer, data); + }, + pong: (ws, data) => { + const peers = getPeers(globalPeers, ws.data.namespace); + const peer = getPeer(ws, peers, baseUtils.sync); + hooks.callHook("pong", peer, data); + }, }, }; }; @@ -156,4 +168,8 @@ class BunPeer extends Peer<{ override terminate(): void { this._internal.ws.terminate(); } + + override ping(data?: unknown): number { + return this._internal.ws.ping(data as any); + } } diff --git a/src/adapters/node.ts b/src/adapters/node.ts index ec1352f..0b1f171 100644 --- a/src/adapters/node.ts +++ b/src/adapters/node.ts @@ -116,6 +116,14 @@ const nodeAdapter: Adapter = (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) => { + hooks.callHook("pong", peer, data); + }); ws.on("error", (error: Error) => { peers.delete(peer); hooks.callHook("error", peer, new WSError(error)); @@ -290,6 +298,10 @@ class NodePeer extends Peer<{ override terminate() { this._internal.ws.terminate(); } + + override ping(data?: unknown): void { + this._internal.ws.ping(data); + } } // --- web compat --- diff --git a/src/adapters/uws.ts b/src/adapters/uws.ts index 30b450a..cd4e01d 100644 --- a/src/adapters/uws.ts +++ b/src/adapters/uws.ts @@ -73,6 +73,18 @@ const uwsAdapter: Adapter = (options = {}) => { const peer = getPeer(ws, peers); 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. + ping(ws, message) { + const peers = getPeers(globalPeers, ws.getUserData().namespace); + const peer = getPeer(ws, peers, baseUtils.sync); + hooks.callHook("ping", peer, new Uint8Array(message)); + }, + pong(ws, message) { + const peers = getPeers(globalPeers, ws.getUserData().namespace); + const peer = getPeer(ws, peers, baseUtils.sync); + hooks.callHook("pong", peer, new Uint8Array(message)); + }, open(ws) { const peers = getPeers(globalPeers, ws.getUserData().namespace); const peer = getPeer(ws, peers, baseUtils.sync); @@ -217,6 +229,10 @@ class UWSPeer extends Peer<{ override terminate(): void { this._internal.uws.close(); } + + override ping(data?: uws.RecognizedString): number { + return this._internal.uws.ping(data); + } } // --- web compat --- diff --git a/src/hooks.ts b/src/hooks.ts index b997b6f..8085312 100644 --- a/src/hooks.ts +++ b/src/hooks.ts @@ -214,4 +214,24 @@ export interface Hooks { /** An error occurs */ error: (peer: Peer, error: WSError) => MaybePromise; + + /** + * 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; + + /** + * 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; } diff --git a/src/peer.ts b/src/peer.ts index 05782cc..4dd28f2 100644 --- a/src/peer.ts +++ b/src/peer.ts @@ -173,6 +173,21 @@ export abstract class Peer { this.close(); } + /** + * Send an application-level WebSocket ping control frame to the client. + * + * Pair with the {@link Hooks.pong} hook (e.g. embedding a timestamp in + * `data`) to measure round-trip latency, or rely on the {@link Hooks.ping} + * hook to observe pings the client sends unprompted. + * + * **Note:** Not all adapters can send a ping frame; unsupported adapters + * warn and no-op. Refer to the + * [compatibility table](https://crossws.h3.dev/guide/peer#compatibility). + */ + ping(_data?: unknown): number | void | undefined { + console.warn("[crossws] `peer.ping()` is not supported by this adapter."); + } + /** Subscribe to a topic */ subscribe(topic: string): void { this._topics.add(topic); diff --git a/src/server/_resolve.ts b/src/server/_resolve.ts index 2f4e1ab..15f48a6 100644 --- a/src/server/_resolve.ts +++ b/src/server/_resolve.ts @@ -3,7 +3,16 @@ import type { Server } from "srvx"; import type { Hooks } from "../hooks"; import type { WSOptions } from "./_types"; -const HOOK_NAMES = ["upgrade", "message", "open", "close", "drain", "error"] as const; +const HOOK_NAMES = [ + "upgrade", + "message", + "open", + "close", + "drain", + "error", + "ping", + "pong", +] as const; // Compile-time guard: if a hook is added to `Hooks` but not listed above, the // leftover key is no longer `never`, so this type resolves to a tuple and the diff --git a/test/adapters/bun.test.ts b/test/adapters/bun.test.ts index daac7a7..ff9c405 100644 --- a/test/adapters/bun.test.ts +++ b/test/adapters/bun.test.ts @@ -1,6 +1,10 @@ import { describe } from "vitest"; import { wsTestsExec } from "../_utils"; +import { wsTests, pingPongTests } from "../tests"; describe("bun", () => { - wsTestsExec("bun run ./bun.ts", { adapter: "bun" }); + wsTestsExec("bun run ./bun.ts", { adapter: "bun" }, (getURL, opts) => { + wsTests(getURL, opts); + pingPongTests(getURL); + }); }); diff --git a/test/adapters/node.test.ts b/test/adapters/node.test.ts index fcbbe0a..445baca 100644 --- a/test/adapters/node.test.ts +++ b/test/adapters/node.test.ts @@ -5,7 +5,7 @@ import { getRandomPort, waitForPort } from "get-port-please"; import nodeAdapter from "../../src/adapters/node"; import { defineHooks } from "../../src/index"; import { createDemo } from "../fixture/_shared"; -import { wsTests } from "../tests"; +import { wsTests, pingPongTests } from "../tests"; import { wsConnect } from "../_utils"; describe("node", () => { @@ -51,6 +51,8 @@ describe("node", () => { adapter: "node", }); + pingPongTests(() => url); + test("forcefully terminates when force=true", async () => { ws.closeAll(undefined, undefined, true); for (const [_ns, peers] of ws.peers) { diff --git a/test/adapters/uws.test.ts b/test/adapters/uws.test.ts index 8ebd677..c7eeea0 100644 --- a/test/adapters/uws.test.ts +++ b/test/adapters/uws.test.ts @@ -9,7 +9,7 @@ import { import uwsAdapter from "../../src/adapters/uws"; import { defineHooks } from "../../src/index"; import { createDemo } from "../fixture/_shared"; -import { wsTests } from "../tests"; +import { wsTests, pingPongTests } from "../tests"; import { wsConnect } from "../_utils"; describe("uws", () => { @@ -67,6 +67,8 @@ describe("uws", () => { wsTests(() => url, { adapter: "uws", }); + + pingPongTests(() => url); }); // Regression: a global `adapter.publish(topic, data)` (no namespace) on a diff --git a/test/fixture/_shared.ts b/test/fixture/_shared.ts index 69e9d13..69f026c 100644 --- a/test/fixture/_shared.ts +++ b/test/fixture/_shared.ts @@ -55,12 +55,22 @@ export function createDemo>( }); break; } + case "ping-me": { + peer.ping("server-ping"); + break; + } default: { peer.send(msgText); peer.publish("chat", msgText); } } }, + ping(peer, data) { + peer.send(`ping-received:${new TextDecoder().decode(data)}`); + }, + pong(peer, data) { + peer.send(`pong-received:${new TextDecoder().decode(data)}`); + }, upgrade(req) { if (req.url.endsWith("?unauthorized")) { throw { diff --git a/test/tests.ts b/test/tests.ts index ec2676c..ddfff7f 100644 --- a/test/tests.ts +++ b/test/tests.ts @@ -1,4 +1,5 @@ import { expect, test } from "vitest"; +import { WebSocket as NodeWebSocket } from "ws"; import { wsConnect } from "./_utils"; export interface WSTestOpts { @@ -199,3 +200,71 @@ export function wsTests(getURL: () => string, opts: WSTestOpts): void { }, ); } + +/** + * Application-level ping/pong control frames aren't reachable through the + * standard `WebSocket` API (browsers/undici auto-answer them invisibly to + * JS), so — unlike {@link wsTests} — this suite connects with the `ws` + * package, which exposes `.ping()`/`.pong()` and the raw `ping`/`pong` + * events. Only wired up for adapters that support it (node, uws, bun); refer + * to the [compatibility table](https://crossws.h3.dev/guide/peer#compatibility). + */ +export function pingPongTests(getURL: () => string): void { + // Queues inbound text messages (attached before `open` resolves, so the + // fixture's immediate "Welcome ..." message can't be missed in the race + // between connecting and a test attaching its own listener) so a test can + // skip past it and assert on the next message deterministically. + const connect = () => { + const client = new NodeWebSocket(getURL()); + const messages: string[] = []; + const waitCallbacks: Record void> = {}; + let nextIndex = 0; + client.on("message", (data) => { + const text = data.toString(); + const index = messages.push(text) - 1; + waitCallbacks[index]?.(text); + delete waitCallbacks[index]; + }); + const next = (): Promise => { + const index = nextIndex++; + if (index < messages.length) { + return Promise.resolve(messages[index]!); + } + return new Promise((resolve) => { + waitCallbacks[index] = resolve; + }); + }; + return new Promise<{ client: NodeWebSocket; next: () => Promise }>( + (resolve, reject) => { + client.once("open", () => resolve({ client, next })); + client.once("error", reject); + }, + ); + }; + + test("ping hook observes an inbound ping from the client", async () => { + const { client, next } = await connect(); + await next(); // "Welcome ..." from the `open` hook + client.ping("client-ping"); + expect(await next()).toBe("ping-received:client-ping"); + client.close(); + }); + + test("pong hook observes an inbound pong from the client", async () => { + const { client, next } = await connect(); + await next(); // "Welcome ..." from the `open` hook + client.pong("client-pong"); + expect(await next()).toBe("pong-received:client-pong"); + client.close(); + }); + + test("peer.ping() sends a ping frame the client receives", async () => { + const { client } = await connect(); + const pingReceived = new Promise((resolve) => { + client.once("ping", (data) => resolve(data.toString())); + }); + client.send("ping-me"); + expect(await pingReceived).toBe("server-ping"); + client.close(); + }); +} From dc8306dfba45f8f2ce59f965b70c3d6607b7d090 Mon Sep 17 00:00:00 2001 From: Pooya Parsa Date: Fri, 3 Jul 2026 10:11:45 +0000 Subject: [PATCH 2/2] fix: surface ping() failures via error hook and correlate keepalive pongs - peer.ping() now catches native validation errors (e.g. non-buffer payloads, >125-byte control frames) and routes them through the error hook instead of crashing the process. - Node's internal keepalive pings are tagged so their echoed pongs don't fire the app-level pong hook; the RTT doc example now embeds a correlatable payload since Bun/uWS keepalive pings can't be tagged. - Warn once per peer (not per call) when ping() is unsupported. - Skip the uWS ping/pong Uint8Array copy when no hook/resolver needs it. - Simplify the pingPongTests message queue to a plain FIFO. Co-Authored-By: Claude Sonnet 5 --- docs/1.guide/3.peer.md | 16 +++++++++++----- src/adapters/bun.ts | 28 ++++++++++++++++++++-------- src/adapters/node.ts | 26 ++++++++++++++++++++++++-- src/adapters/uws.ts | 39 +++++++++++++++++++++++++++++++-------- src/peer.ts | 14 ++++++++++++-- test/tests.ts | 23 +++++++++++++---------- 6 files changed, 111 insertions(+), 35 deletions(-) diff --git a/docs/1.guide/3.peer.md b/docs/1.guide/3.peer.md index 959822b..443f14c 100644 --- a/docs/1.guide/3.peer.md +++ b/docs/1.guide/3.peer.md @@ -137,7 +137,7 @@ To gracefully close the connection, use `peer.close()`. Send an application-level WebSocket ping control frame to the client. -Pair with the [`pong` hook](/guide/hooks) — e.g. embedding a timestamp in `data` — to measure round-trip latency: +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"; @@ -145,12 +145,15 @@ import { defineHooks } from "crossws"; const hooks = defineHooks({ message(peer, message) { if (message.text() === "measure-latency") { - peer.context.pingSentAt = Date.now(); - peer.ping(); + // 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) { - const rtt = Date.now() - (peer.context.pingSentAt as number); + 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`); }, }); @@ -158,6 +161,9 @@ const hooks = defineHooks({ 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). diff --git a/src/adapters/bun.ts b/src/adapters/bun.ts index 2f86c02..067670d 100644 --- a/src/adapters/bun.ts +++ b/src/adapters/bun.ts @@ -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"; @@ -65,36 +66,36 @@ const bunAdapter: Adapter = (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); + 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); + const peer = getPeer(ws, peers, baseUtils.sync, hooks); hooks.callHook("pong", peer, data); }, }, @@ -108,7 +109,8 @@ export default bunAdapter; function getPeer( ws: ServerWebSocket, peers: Set, - sync?: SyncDriver, + sync: SyncDriver | undefined, + hooks: AdapterHookable, ): BunPeer { if (ws.data.peer) { return ws.data.peer; @@ -119,6 +121,7 @@ function getPeer( peers, namespace: ws.data.namespace, sync, + hooks, }); ws.data.peer = peer; return peer; @@ -130,6 +133,7 @@ class BunPeer extends Peer<{ request: Request; peers: Set; sync?: SyncDriver; + hooks: AdapterHookable; }> { override get remoteAddress(): string { return this._internal.ws.remoteAddress; @@ -170,6 +174,14 @@ class BunPeer extends Peer<{ } override ping(data?: unknown): number { - return this._internal.ws.ping(data as any); + // 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; + } } } diff --git a/src/adapters/node.ts b/src/adapters/node.ts index 0b1f171..9c95944 100644 --- a/src/adapters/node.ts +++ b/src/adapters/node.ts @@ -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, @@ -89,6 +95,7 @@ const nodeAdapter: Adapter = (options = {}) => { nodeReq, namespace: nodeReq._namespace, sync: baseUtils.sync, + hooks, }); peers.add(peer); liveSockets.add(ws as HeartbeatWS); @@ -122,6 +129,12 @@ const nodeAdapter: Adapter = (options = {}) => { 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) => { @@ -168,7 +181,7 @@ const nodeAdapter: Adapter = (options = {}) => { } ws._isAlive = false; try { - ws.ping(); + ws.ping(HEARTBEAT_PING); } catch { // socket may have raced into CLOSING between the sweep and the ping } @@ -253,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; @@ -300,7 +314,15 @@ class NodePeer extends Peer<{ } override ping(data?: unknown): void { - this._internal.ws.ping(data); + // `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)); + } } } diff --git a/src/adapters/uws.ts b/src/adapters/uws.ts index cd4e01d..9b146e7 100644 --- a/src/adapters/uws.ts +++ b/src/adapters/uws.ts @@ -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"; @@ -54,7 +55,7 @@ const uwsAdapter: Adapter = (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, { @@ -65,29 +66,36 @@ const uwsAdapter: Adapter = (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); + 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); + 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); }, @@ -159,7 +167,12 @@ export default uwsAdapter; // --- peer --- -function getPeer(uws: uws.WebSocket, peers: Set, sync?: SyncDriver): UWSPeer { +function getPeer( + uws: uws.WebSocket, + peers: Set, + sync: SyncDriver | undefined, + hooks: AdapterHookable, +): UWSPeer { const uwsData = uws.getUserData(); if (uwsData.peer) { return uwsData.peer; @@ -172,6 +185,7 @@ function getPeer(uws: uws.WebSocket, peers: Set, sync?: SyncD namespace: uwsData.namespace, uwsData, sync, + hooks, }); uwsData.peer = peer; return peer; @@ -185,6 +199,7 @@ class UWSPeer extends Peer<{ ws: UwsWebSocketProxy; uwsData: UserData; sync?: SyncDriver; + hooks: AdapterHookable; }> { override get remoteAddress(): string | undefined { try { @@ -231,7 +246,15 @@ class UWSPeer extends Peer<{ } override ping(data?: uws.RecognizedString): number { - return this._internal.uws.ping(data); + // 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; + } } } diff --git a/src/peer.ts b/src/peer.ts index 4dd28f2..912ebdd 100644 --- a/src/peer.ts +++ b/src/peer.ts @@ -42,6 +42,7 @@ export abstract class Peer { protected _id?: string; #ws?: Partial; + #pingUnsupportedWarned = false; constructor(internal: Internal) { this._topics = new Set(); @@ -180,12 +181,21 @@ export abstract class Peer { * `data`) to measure round-trip latency, or rely on the {@link Hooks.ping} * hook to observe pings the client sends unprompted. * + * `data` is optional and, per RFC 6455, must not exceed 125 bytes (a ping is + * a control frame); an over-long or otherwise invalid payload is reported via + * the {@link Hooks.error} hook instead of being sent. + * * **Note:** Not all adapters can send a ping frame; unsupported adapters - * warn and no-op. Refer to the + * warn once and no-op. Refer to the * [compatibility table](https://crossws.h3.dev/guide/peer#compatibility). */ ping(_data?: unknown): number | void | undefined { - console.warn("[crossws] `peer.ping()` is not supported by this adapter."); + // Warn once per peer instead of on every call, so a heartbeat loop against + // an unsupported adapter can't spam the console. + if (!this.#pingUnsupportedWarned) { + this.#pingUnsupportedWarned = true; + console.warn("[crossws] `peer.ping()` is not supported by this adapter."); + } } /** Subscribe to a topic */ diff --git a/test/tests.ts b/test/tests.ts index ddfff7f..d59b79b 100644 --- a/test/tests.ts +++ b/test/tests.ts @@ -216,22 +216,25 @@ export function pingPongTests(getURL: () => string): void { // skip past it and assert on the next message deterministically. const connect = () => { const client = new NodeWebSocket(getURL()); - const messages: string[] = []; - const waitCallbacks: Record void> = {}; - let nextIndex = 0; + const queue: string[] = []; + let pending: ((message: string) => void) | undefined; client.on("message", (data) => { const text = data.toString(); - const index = messages.push(text) - 1; - waitCallbacks[index]?.(text); - delete waitCallbacks[index]; + if (pending) { + pending(text); + pending = undefined; + } else { + queue.push(text); + } }); + // The tests only ever await `next()` sequentially, so a plain FIFO queue + // with a single pending resolver is enough. const next = (): Promise => { - const index = nextIndex++; - if (index < messages.length) { - return Promise.resolve(messages[index]!); + if (queue.length > 0) { + return Promise.resolve(queue.shift()!); } return new Promise((resolve) => { - waitCallbacks[index] = resolve; + pending = resolve; }); }; return new Promise<{ client: NodeWebSocket; next: () => Promise }>(