diff --git a/extensions/browser/src/browser/cdp.close-tracked.test.ts b/extensions/browser/src/browser/cdp.close-tracked.test.ts index 9f7c0508d474..634fd5908e2a 100644 --- a/extensions/browser/src/browser/cdp.close-tracked.test.ts +++ b/extensions/browser/src/browser/cdp.close-tracked.test.ts @@ -1,8 +1,8 @@ import { createServer } from "node:http"; import type { AddressInfo } from "node:net"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import { afterEach, describe, expect, it } from "vitest"; import { type WebSocket, WebSocketServer } from "ws"; -import { rawDataToString } from "../infra/ws.js"; import "../test-support/browser-security.mock.js"; import { closeTrackedCdpTarget, resolveCdpTabOwnership } from "./cdp.helpers.js"; diff --git a/extensions/browser/src/browser/cdp.helpers.internal.test.ts b/extensions/browser/src/browser/cdp.helpers.internal.test.ts index 415187bf9092..0596a9a6cb2d 100644 --- a/extensions/browser/src/browser/cdp.helpers.internal.test.ts +++ b/extensions/browser/src/browser/cdp.helpers.internal.test.ts @@ -1,10 +1,10 @@ // Browser tests cover cdp.helpers.internal plugin behavior. import { createServer } from "node:http"; import type { Socket } from "node:net"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { WebSocketServer } from "ws"; import { toErrorObject } from "../infra/errors.js"; -import { rawDataToString } from "../infra/ws.js"; const fetchWithSsrFGuardMock = vi.hoisted(() => vi.fn()); const sleepWithAbortMock = vi.hoisted(() => diff --git a/extensions/browser/src/browser/cdp.helpers.ts b/extensions/browser/src/browser/cdp.helpers.ts index 16bf084dd0ab..e36e3a26bf05 100644 --- a/extensions/browser/src/browser/cdp.helpers.ts +++ b/extensions/browser/src/browser/cdp.helpers.ts @@ -9,6 +9,7 @@ import { parseBrowserHttpUrl, redactCdpUrl } from "openclaw/plugin-sdk/browser-c import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; import { sleepWithAbort } from "openclaw/plugin-sdk/runtime-env"; import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import WebSocket from "ws"; import { isLoopbackHost } from "../gateway/net.js"; import { @@ -16,7 +17,6 @@ import { type SsrFPolicy, resolvePinnedHostnameWithPolicy, } from "../infra/net/ssrf.js"; -import { rawDataToString } from "../infra/ws.js"; import { redactToolPayloadText } from "../logging/redact.js"; import { getDirectAgentForCdp, diff --git a/extensions/browser/src/browser/cdp.internal.test.ts b/extensions/browser/src/browser/cdp.internal.test.ts index 0b27837dbccb..ae591af5dbac 100644 --- a/extensions/browser/src/browser/cdp.internal.test.ts +++ b/extensions/browser/src/browser/cdp.internal.test.ts @@ -1,7 +1,7 @@ +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; // Browser tests cover cdp.internal plugin behavior. import { afterEach, describe, expect, it } from "vitest"; import { WebSocketServer } from "ws"; -import { rawDataToString } from "../infra/ws.js"; import "../test-support/browser-security.mock.js"; import { type AriaSnapshotNode, diff --git a/extensions/browser/src/browser/cdp.test.ts b/extensions/browser/src/browser/cdp.test.ts index 0cd0964f95ec..ec968c496865 100644 --- a/extensions/browser/src/browser/cdp.test.ts +++ b/extensions/browser/src/browser/cdp.test.ts @@ -2,10 +2,10 @@ import { createServer } from "node:http"; import type { AddressInfo } from "node:net"; import type { Duplex } from "node:stream"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { type WebSocket, WebSocketServer } from "ws"; import { SsrFBlockedError } from "../infra/net/ssrf.js"; -import { rawDataToString } from "../infra/ws.js"; import "../test-support/browser-security.mock.js"; import { closeTrackedCdpTarget, diff --git a/extensions/browser/src/browser/chrome.diagnostics.ts b/extensions/browser/src/browser/chrome.diagnostics.ts index d2e6c9571eb4..654fa25fd533 100644 --- a/extensions/browser/src/browser/chrome.diagnostics.ts +++ b/extensions/browser/src/browser/chrome.diagnostics.ts @@ -6,8 +6,8 @@ import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; * and formats status output for browser doctor/status flows. */ import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import type { SsrFPolicy } from "../infra/net/ssrf.js"; -import { rawDataToString } from "../infra/ws.js"; import { redactSensitiveText } from "../logging/redact.js"; import { CHROME_REACHABILITY_TIMEOUT_MS, CHROME_WS_READY_TIMEOUT_MS } from "./cdp-timeouts.js"; import { diff --git a/extensions/browser/src/browser/chrome.internal.test.ts b/extensions/browser/src/browser/chrome.internal.test.ts index e07ede594410..8bc028deeeb5 100644 --- a/extensions/browser/src/browser/chrome.internal.test.ts +++ b/extensions/browser/src/browser/chrome.internal.test.ts @@ -7,9 +7,9 @@ import type { AddressInfo } from "node:net"; import os from "node:os"; import path from "node:path"; import { createOpenClawTestState } from "openclaw/plugin-sdk/test-state"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { WebSocketServer } from "ws"; -import { rawDataToString } from "../infra/ws.js"; const spawnMock = vi.hoisted(() => vi.fn()); diff --git a/extensions/browser/src/browser/chrome.test.ts b/extensions/browser/src/browser/chrome.test.ts index 94a20b6173f8..26f57cbb6329 100644 --- a/extensions/browser/src/browser/chrome.test.ts +++ b/extensions/browser/src/browser/chrome.test.ts @@ -4,9 +4,9 @@ import fs from "node:fs"; import { createServer } from "node:http"; import { createServer as createTcpServer } from "node:net"; import type { AddressInfo } from "node:net"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { WebSocketServer } from "ws"; -import { rawDataToString } from "../infra/ws.js"; import { diagnoseChromeCdp, formatChromeCdpDiagnostic } from "./chrome.diagnostics.js"; import { parseBrowserMajorVersion, diff --git a/extensions/browser/src/browser/extension-relay/relay-server.ts b/extensions/browser/src/browser/extension-relay/relay-server.ts index 5a9f19df08a1..651d4f067182 100644 --- a/extensions/browser/src/browser/extension-relay/relay-server.ts +++ b/extensions/browser/src/browser/extension-relay/relay-server.ts @@ -3,9 +3,9 @@ import crypto from "node:crypto"; import http, { type IncomingMessage, type Server, type ServerResponse } from "node:http"; import type { Duplex } from "node:stream"; import { safeEqualSecret } from "openclaw/plugin-sdk/security-runtime"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import { WebSocketServer, type RawData, type WebSocket } from "ws"; import { isLoopbackHost } from "../../gateway/net.js"; -import { rawDataToString } from "../../infra/ws.js"; import { createSubsystemLogger } from "../../logging/subsystem.js"; import { BROWSER_RELAY_AUTH_CHALLENGE_PATH, diff --git a/extensions/browser/src/infra/ws.ts b/extensions/browser/src/infra/ws.ts deleted file mode 100644 index 5216d411b4b0..000000000000 --- a/extensions/browser/src/infra/ws.ts +++ /dev/null @@ -1,22 +0,0 @@ -/** - * WebSocket payload normalization helpers for Browser gateway transports. - */ -/** Converts raw WebSocket payload shapes into UTF-8 strings. */ -export function rawDataToString(data: unknown): string { - if (typeof data === "string") { - return data; - } - if (Buffer.isBuffer(data)) { - return data.toString("utf8"); - } - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - if (ArrayBuffer.isView(data)) { - return Buffer.from(data.buffer, data.byteOffset, data.byteLength).toString("utf8"); - } - if (data instanceof ArrayBuffer) { - return Buffer.from(data).toString("utf8"); - } - return String(data); -} diff --git a/extensions/clickclack/src/gateway.ts b/extensions/clickclack/src/gateway.ts index 3616c0f7031f..7220afc58047 100644 --- a/extensions/clickclack/src/gateway.ts +++ b/extensions/clickclack/src/gateway.ts @@ -6,6 +6,7 @@ import type { ChannelGatewayContext } from "openclaw/plugin-sdk/channel-contract import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime"; import { channelReadyPatch, channelStoppedPatch } from "openclaw/plugin-sdk/gateway-runtime"; import { sleepWithAbort } from "openclaw/plugin-sdk/runtime-env"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import type { RawData } from "ws"; import { resolveClickClackInboundAccess } from "./access.js"; import { resolveClickClackAccount } from "./accounts.js"; @@ -53,22 +54,9 @@ async function resolveEventMessage(params: { } } -function decodeSocketMessage(data: RawData): string { - if (typeof data === "string") { - return data; - } - if (Buffer.isBuffer(data)) { - return data.toString("utf8"); - } - if (data instanceof ArrayBuffer) { - return Buffer.from(data).toString("utf8"); - } - return Buffer.concat(data).toString("utf8"); -} - function parseSocketEvent(data: RawData): ClickClackEvent | null { try { - return JSON.parse(decodeSocketMessage(data)) as ClickClackEvent; + return JSON.parse(rawDataToString(data)) as ClickClackEvent; } catch { return null; } diff --git a/extensions/openai/realtime-quicksilver-bridge.ts b/extensions/openai/realtime-quicksilver-bridge.ts index 393b7030c04b..d58e007f3bd9 100644 --- a/extensions/openai/realtime-quicksilver-bridge.ts +++ b/extensions/openai/realtime-quicksilver-bridge.ts @@ -17,6 +17,7 @@ import { type RealtimeVoiceSessionConnection, type RealtimeVoiceToolResultOptions, } from "openclaw/plugin-sdk/realtime-voice"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import WebSocket, { type RawData } from "ws"; import { connectOpenAIQuicksilverSideband, @@ -46,16 +47,6 @@ type OpenAIQuicksilverVoiceBridgeConfig = RealtimeVoiceBridgeCreateRequest & { webSocketFactory?: OpenAIQuicksilverSocketFactory; }; -function decodeTextFrame(data: RawData): string { - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - if (data instanceof ArrayBuffer) { - return Buffer.from(data).toString("utf8"); - } - return data.toString("utf8"); -} - function toolResultText(result: unknown): string { if (typeof result === "string") { return result; @@ -212,7 +203,7 @@ export class OpenAIQuicksilverVoiceBridge implements RealtimeVoiceBridge { } return; } - const payload = decodeTextFrame(data); + const payload = rawDataToString(data); captureWsEvent({ url, direction: "inbound", @@ -265,7 +256,7 @@ export class OpenAIQuicksilverVoiceBridge implements RealtimeVoiceBridge { ); for (const frame of connected.bufferedFrames) { if (!frame.isBinary) { - const event = parseOpenAIQuicksilverEvent(decodeTextFrame(frame.data)); + const event = parseOpenAIQuicksilverEvent(rawDataToString(frame.data)); if (event) { this.handleEvent(event, connection, settleReady, failStartup); } diff --git a/extensions/openai/realtime-quicksilver-delegation-controller.ts b/extensions/openai/realtime-quicksilver-delegation-controller.ts index 00e91301e0be..26e9d0d07895 100644 --- a/extensions/openai/realtime-quicksilver-delegation-controller.ts +++ b/extensions/openai/realtime-quicksilver-delegation-controller.ts @@ -1,5 +1,6 @@ import type { PluginLogger } from "openclaw/plugin-sdk/plugin-entry"; import type { RealtimeVoiceAgentConsultRunner } from "openclaw/plugin-sdk/realtime-voice"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import type { RawData } from "ws"; import { buildOpenAIQuicksilverDelegationPrompt, @@ -43,16 +44,6 @@ function shortFailureReason(error: unknown): string { return toError(error).message.replaceAll(/\s+/g, " ").trim().slice(0, 180) || "unknown error"; } -function decodeTextFrame(data: RawData): string { - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - if (data instanceof ArrayBuffer) { - return Buffer.from(data).toString("utf8"); - } - return data.toString("utf8"); -} - function readWireEventType(payload: string): string | undefined { try { const decoded = JSON.parse(payload) as Record; @@ -78,7 +69,7 @@ export class OpenAIQuicksilverDelegationController { this.fail(new Error("OpenAI GPT-Live sideband returned an unexpected binary frame")); return; } - const payload = decodeTextFrame(data); + const payload = rawDataToString(data); if (this.options.onWireEventType) { const eventType = readWireEventType(payload); if (eventType) { diff --git a/extensions/slack/src/monitor/relay-source.ts b/extensions/slack/src/monitor/relay-source.ts index d0df06d9c528..d4a813e9b3f6 100644 --- a/extensions/slack/src/monitor/relay-source.ts +++ b/extensions/slack/src/monitor/relay-source.ts @@ -8,6 +8,7 @@ import { type RuntimeEnv, } from "openclaw/plugin-sdk/runtime-env"; import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; import WebSocket, { type ClientOptions, type RawData } from "ws"; import type { SlackSendIdentity } from "../send.js"; import type { SlackMessageEvent } from "../types.js"; @@ -288,19 +289,6 @@ export function parseRelayFrame(data: RawData): unknown { } } -function rawDataToString(data: RawData): string { - if (typeof data === "string") { - return data; - } - if (Buffer.isBuffer(data)) { - return data.toString("utf8"); - } - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - return Buffer.from(data).toString("utf8"); -} - function extractRelaySlackMessageEvent( frame: unknown, ): { deliveryId: string; message: SlackMessageEvent; route: SlackRelayRoute } | undefined { diff --git a/extensions/xai/tts.ts b/extensions/xai/tts.ts index 6f59f83453c7..1c2d55c3dcd1 100644 --- a/extensions/xai/tts.ts +++ b/extensions/xai/tts.ts @@ -11,7 +11,8 @@ import { fetchWithSsrFGuard, ssrfPolicyFromHttpBaseUrlAllowedHostname, } from "openclaw/plugin-sdk/ssrf-runtime"; -import WebSocket, { type RawData } from "ws"; +import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress"; +import WebSocket from "ws"; import { XAI_BASE_URL } from "./model-definitions.js"; import { isValidXaiTtsVoice, @@ -134,22 +135,6 @@ function assertXaiNativeTtsStreamEndpoint(baseUrl: string): void { } } -function decodeWebSocketTextMessage(data: RawData): string { - if (typeof data === "string") { - return data; - } - if (Buffer.isBuffer(data)) { - return data.toString("utf8"); - } - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - if (data instanceof ArrayBuffer) { - return Buffer.from(data).toString("utf8"); - } - throw new Error("xAI TTS stream received unsupported WebSocket message payload"); -} - export async function xaiTTSStream(params: { text: string; apiKey: string; @@ -384,7 +369,7 @@ export async function xaiTTSStream(params: { return; } try { - const payload = decodeWebSocketTextMessage(data); + const payload = rawDataToString(data); handleServerEvent(JSON.parse(payload) as XaiTtsStreamServerEvent); } catch (error) { failStream(error instanceof Error ? error : new Error(String(error))); diff --git a/packages/gateway-client/README.md b/packages/gateway-client/README.md index 7ac8d54a50ff..8a4440033928 100644 --- a/packages/gateway-client/README.md +++ b/packages/gateway-client/README.md @@ -44,6 +44,8 @@ client surface. until the event loop can process Gateway IO. - `@openclaw/gateway-client/timeouts` exports timeout constants and safe timer resolution helpers. +- `@openclaw/gateway-client/websocket-data` converts every Node `ws` raw-data + shape to UTF-8 text. ## Node quickstart diff --git a/packages/gateway-client/package.json b/packages/gateway-client/package.json index 1f5b130e7143..83c951be3db4 100644 --- a/packages/gateway-client/package.json +++ b/packages/gateway-client/package.json @@ -49,10 +49,15 @@ "types": "./dist/timeouts.d.mts", "import": "./dist/timeouts.mjs", "default": "./dist/timeouts.mjs" + }, + "./websocket-data": { + "types": "./dist/websocket-data.d.mts", + "import": "./dist/websocket-data.mjs", + "default": "./dist/websocket-data.mjs" } }, "scripts": { - "build": "tsdown src/index.ts src/browser.ts src/readiness.ts src/timeouts.ts --no-config --platform node --format esm --dts --out-dir dist --clean" + "build": "tsdown src/index.ts src/browser.ts src/readiness.ts src/timeouts.ts src/websocket-data.ts --no-config --platform node --format esm --dts --out-dir dist --clean" }, "dependencies": { "@openclaw/gateway-protocol": "workspace:*", diff --git a/packages/gateway-client/src/client.watchdog.test.ts b/packages/gateway-client/src/client.watchdog.test.ts index 5a14647d5f6f..5a47c0fbf434 100644 --- a/packages/gateway-client/src/client.watchdog.test.ts +++ b/packages/gateway-client/src/client.watchdog.test.ts @@ -43,6 +43,7 @@ test("decodes every ws raw-data shape", () => { expect(rawDataToString(Buffer.from("buffer"))).toBe("buffer"); expect(rawDataToString(Uint8Array.from(Buffer.from("array-buffer")).buffer)).toBe("array-buffer"); expect(rawDataToString([Buffer.from("frag"), Buffer.from("ments")])).toBe("fragments"); + expect(rawDataToString(Buffer.from([0xe9]), "latin1")).toBe("é"); }); type ProtocolHarness = { diff --git a/packages/gateway-client/src/websocket-data.ts b/packages/gateway-client/src/websocket-data.ts index 62d73056468c..e0df7b58e88f 100644 --- a/packages/gateway-client/src/websocket-data.ts +++ b/packages/gateway-client/src/websocket-data.ts @@ -2,9 +2,11 @@ import { Buffer } from "node:buffer"; import type { RawData } from "ws"; -export function rawDataToString(data: RawData): string { +export function rawDataToString(data: RawData, encoding: BufferEncoding = "utf8"): string { if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); + return Buffer.concat(data).toString(encoding); } - return data instanceof ArrayBuffer ? Buffer.from(data).toString("utf8") : data.toString("utf8"); + return data instanceof ArrayBuffer + ? Buffer.from(data).toString(encoding) + : data.toString(encoding); } diff --git a/packages/sdk/src/app-sdk-composed-resources.e2e.test.ts b/packages/sdk/src/app-sdk-composed-resources.e2e.test.ts index 51077181f2b7..16ece49b2bed 100644 --- a/packages/sdk/src/app-sdk-composed-resources.e2e.test.ts +++ b/packages/sdk/src/app-sdk-composed-resources.e2e.test.ts @@ -3,6 +3,7 @@ import fs from "node:fs/promises"; import type { AddressInfo } from "node:net"; import os from "node:os"; import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import type { TSchema } from "typebox"; import { Value } from "typebox/value"; import { afterEach, describe, expect, it, vi } from "vitest"; @@ -40,7 +41,6 @@ import { } from "../../../src/gateway/test-helpers.js"; import { emitAgentEvent } from "../../../src/infra/agent-events.js"; import { registerAgentRunContext } from "../../../src/infra/agent-run-registry.js"; -import { rawDataToString } from "../../../src/infra/ws.js"; import { withTimeout } from "../../../src/utils/with-timeout.js"; import { GatewayClientTransport, OpenClaw, type OpenClawEvent } from "./index.js"; diff --git a/packages/sdk/src/index.e2e.test.ts b/packages/sdk/src/index.e2e.test.ts index 9c40999f8d05..efef9bbab98e 100644 --- a/packages/sdk/src/index.e2e.test.ts +++ b/packages/sdk/src/index.e2e.test.ts @@ -1,12 +1,12 @@ // OpenClaw SDK tests cover index behavior. import type { AddressInfo } from "node:net"; import net from "node:net"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { afterEach, describe, expect, it } from "vitest"; import { WebSocketServer, type WebSocket } from "ws"; import { installGatewayTestHooks, startServer } from "../../../src/gateway/test-helpers.js"; import { emitAgentEvent } from "../../../src/infra/agent-events.js"; import { registerAgentRunContext } from "../../../src/infra/agent-run-registry.js"; -import { rawDataToString } from "../../../src/infra/ws.js"; import { withTimeout } from "../../../src/utils/with-timeout.js"; import { GatewayClientTransport, OpenClaw } from "./index.js"; diff --git a/scripts/lib/gateway-ws-client.ts b/scripts/lib/gateway-ws-client.ts index b573ef45db0f..16cedbdb6b98 100644 --- a/scripts/lib/gateway-ws-client.ts +++ b/scripts/lib/gateway-ws-client.ts @@ -1,16 +1,7 @@ // Gateway Ws Client script supports OpenClaw repository automation. -import { Buffer } from "node:buffer"; import { randomUUID } from "node:crypto"; import WebSocket from "ws"; - -// Release Docker images ship this script without src/, so keep transport -// normalization self-contained instead of importing a core-only helper. -export function rawDataToString(data: WebSocket.RawData): string { - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - return data instanceof ArrayBuffer ? Buffer.from(data).toString("utf8") : data.toString("utf8"); -} +import { rawDataToString } from "../../packages/gateway-client/src/websocket-data.ts"; type GatewayReqFrame = { type: "req"; id: string; method: string; params?: unknown }; type GatewayResFrame = { diff --git a/src/gateway/minimal-gateway.test-helpers.ts b/src/gateway/minimal-gateway.test-helpers.ts index dce123d27f20..eeaf28b7582a 100644 --- a/src/gateway/minimal-gateway.test-helpers.ts +++ b/src/gateway/minimal-gateway.test-helpers.ts @@ -1,9 +1,9 @@ // Minimal Gateway websocket test helpers. // Provides small fake-server frames plus an isolated real-Gateway boundary harness. import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { WebSocket, type WebSocketServer } from "ws"; import { PROTOCOL_VERSION } from "../../packages/gateway-protocol/src/index.js"; -import { rawDataToString } from "../infra/ws.js"; import { toAgentRequestSessionKey } from "../routing/session-key.js"; import { getFreePort } from "../test-utils/ports.js"; diff --git a/src/gateway/node-registry.ws-lifecycle.test.ts b/src/gateway/node-registry.ws-lifecycle.test.ts index ca9a796c8207..241c6bffa8cd 100644 --- a/src/gateway/node-registry.ws-lifecycle.test.ts +++ b/src/gateway/node-registry.ws-lifecycle.test.ts @@ -1,9 +1,9 @@ import { once } from "node:events"; import type { AddressInfo } from "node:net"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { afterEach, describe, expect, it, vi } from "vitest"; import { WebSocket, WebSocketServer, type RawData } from "ws"; import { setActiveNodeContext } from "../infra/active-node-context.js"; -import { rawDataToString } from "../infra/ws.js"; import { NodeRegistry } from "./node-registry.js"; import type { GatewayWsClient } from "./server/ws-types.js"; diff --git a/src/gateway/server.agent.rpc-contracts.test.ts b/src/gateway/server.agent.rpc-contracts.test.ts index 9a7880c1ebe4..4d6c48e8da78 100644 --- a/src/gateway/server.agent.rpc-contracts.test.ts +++ b/src/gateway/server.agent.rpc-contracts.test.ts @@ -1,8 +1,8 @@ +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; // Real Gateway WebSocket proof for agent delivery fallback, response ordering, and idempotency. import { afterAll, beforeAll, beforeEach, describe, expect, test, vi } from "vitest"; import type { RawData, WebSocket } from "ws"; import { createDeferred } from "../../test/helpers/promise.js"; -import { rawDataToString } from "../infra/ws.js"; import { startGatewayServerHarness, type GatewayServerHarness } from "./server.e2e-ws-harness.js"; import { agentCommand, installGatewayTestHooks, onceMessage } from "./test-helpers.js"; diff --git a/src/gateway/server.chat-abort-dispatch-rejection.test.ts b/src/gateway/server.chat-abort-dispatch-rejection.test.ts index a1681438db27..0a77836dd82a 100644 --- a/src/gateway/server.chat-abort-dispatch-rejection.test.ts +++ b/src/gateway/server.chat-abort-dispatch-rejection.test.ts @@ -1,12 +1,12 @@ // Real WebSocket coverage for abort ownership when an in-flight dispatch rejects. import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { afterAll, afterEach, beforeAll, describe, expect, test, vi } from "vitest"; import { createDeferred } from "../../test/helpers/promise.js"; import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; import type { GetReplyOptions } from "../auto-reply/get-reply-options.types.js"; import { clearConfigCache } from "../config/config.js"; import { emitAgentEvent } from "../infra/agent-events.js"; -import { rawDataToString } from "../infra/ws.js"; import { interruptSessionWorkAdmissions } from "../sessions/session-lifecycle-admission.js"; import { connectOk, diff --git a/src/gateway/server/ws-connection/message-handler.ts b/src/gateway/server/ws-connection/message-handler.ts index 136050635e15..82b9f9e6572e 100644 --- a/src/gateway/server/ws-connection/message-handler.ts +++ b/src/gateway/server/ws-connection/message-handler.ts @@ -1,4 +1,5 @@ // WebSocket message handler validates frames, dispatches gateway RPCs, manages pairing, and reports responses. +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import type { RawData } from "ws"; import { GATEWAY_CLIENT_IDS, @@ -22,7 +23,7 @@ import { createDiagnosticTraceContext, runWithDiagnosticTraceContext, } from "../../../infra/diagnostic-trace-context.js"; -import { rawDataByteLength, rawDataToString } from "../../../infra/ws.js"; +import { rawDataByteLength } from "../../../infra/ws.js"; import { logRejectedLargePayload } from "../../../logging/diagnostic-payload.js"; import { getGatewaySuspendAdmissionPhase, diff --git a/src/gateway/server/ws-connection/worker-connection.ts b/src/gateway/server/ws-connection/worker-connection.ts index 89f9a441f6d6..870682f97523 100644 --- a/src/gateway/server/ws-connection/worker-connection.ts +++ b/src/gateway/server/ws-connection/worker-connection.ts @@ -1,3 +1,4 @@ +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import type { RawData, WebSocket } from "ws"; import { ErrorCodes, @@ -45,7 +46,7 @@ import { validateWorkerInferenceStartParams, } from "../../../../packages/gateway-protocol/src/schema/worker-inference.js"; import { GATEWAY_STARTUP_RETRY_AFTER_MS } from "../../../../packages/gateway-protocol/src/startup-unavailable.js"; -import { rawDataByteLength, rawDataToString } from "../../../infra/ws.js"; +import { rawDataByteLength } from "../../../infra/ws.js"; import { tryBeginGatewayRootWorkAdmission } from "../../../process/gateway-work-admission.js"; import type { WorkerConnectionIdentity } from "../../worker-environments/connection-identity.js"; import type { GatewayWsClient, WsHandshakePhase } from "../ws-types.js"; diff --git a/src/gateway/session-message-events.test.ts b/src/gateway/session-message-events.test.ts index 574388466e8b..2411ab19c27b 100644 --- a/src/gateway/session-message-events.test.ts +++ b/src/gateway/session-message-events.test.ts @@ -1,6 +1,7 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; /** * Session message event indexing and broadcast tests. */ @@ -25,7 +26,6 @@ import { appendAssistantMessageToSessionTranscript } from "../config/sessions/tr import type { OpenClawConfig } from "../config/types.openclaw.js"; import { emitAgentEvent } from "../infra/agent-events.js"; import { claimAgentRunContext, clearAgentRunContext } from "../infra/agent-run-registry.js"; -import { rawDataToString } from "../infra/ws.js"; import { emitSessionLifecycleEvent } from "../sessions/session-lifecycle-events.js"; import * as transcriptEvents from "../sessions/transcript-events.js"; import { emitSessionTranscriptUpdate } from "../sessions/transcript-events.js"; diff --git a/src/gateway/test-helpers.e2e.ts b/src/gateway/test-helpers.e2e.ts index 827a8776b649..b3dc477b05f3 100644 --- a/src/gateway/test-helpers.e2e.ts +++ b/src/gateway/test-helpers.e2e.ts @@ -3,6 +3,7 @@ import { writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { normalizeLowercaseStringOrEmpty } from "@openclaw/normalization-core/string-coerce"; import { WebSocket } from "ws"; import { type HelloOk, PROTOCOL_VERSION } from "../../packages/gateway-protocol/src/index.js"; @@ -14,7 +15,6 @@ import { publicKeyRawBase64UrlFromPem, signDevicePayload, } from "../infra/device-identity.js"; -import { rawDataToString } from "../infra/ws.js"; import { captureEnv } from "../test-utils/env.js"; import { getDeterministicFreePortBlock } from "../test-utils/ports.js"; import { diff --git a/src/gateway/test-helpers.server.ts b/src/gateway/test-helpers.server.ts index b221386b9d73..ac8ace2085d8 100644 --- a/src/gateway/test-helpers.server.ts +++ b/src/gateway/test-helpers.server.ts @@ -3,9 +3,10 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { normalizeLowercaseStringOrEmpty } from "@openclaw/normalization-core/string-coerce"; -import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import "./test-helpers.mocks.js"; +import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import { afterAll, afterEach, beforeAll, beforeEach, expect, vi } from "vitest"; import { WebSocket } from "ws"; import { PROTOCOL_VERSION } from "../../packages/gateway-protocol/src/index.js"; @@ -37,7 +38,6 @@ import { } from "../infra/restart.js"; import { normalizeLegacySessionEntryDelivery } from "../infra/state-migrations.legacy-session-store.js"; import { drainSystemEvents, peekSystemEvents } from "../infra/system-events.js"; -import { rawDataToString } from "../infra/ws.js"; import { resetLogger, setLoggerOverride } from "../logging.js"; import type { ChannelRouteRef } from "../plugin-sdk/channel-route.js"; import { resetGatewayWorkAdmission } from "../process/gateway-work-admission.js"; diff --git a/src/infra/ws.test.ts b/src/infra/ws.test.ts index 26c3e6c568f7..e9e27d76d572 100644 --- a/src/infra/ws.test.ts +++ b/src/infra/ws.test.ts @@ -1,15 +1,14 @@ -// Covers WebSocket raw payload decoding. +// Covers WebSocket raw payload sizing. import { Buffer } from "node:buffer"; import { describe, expect, it } from "vitest"; -import { rawDataByteLength, rawDataToString } from "./ws.js"; +import { rawDataByteLength } from "./ws.js"; -describe("WebSocket raw data", () => { +describe("rawDataByteLength", () => { it.each([ ["Buffer", Buffer.from("hello")], ["Buffer[]", [Buffer.from("he"), Buffer.from("llo")]], ["ArrayBuffer", Uint8Array.from([104, 101, 108, 108, 111]).buffer], ])("handles %s", (_name, data) => { - expect(rawDataToString(data)).toBe("hello"); expect(rawDataByteLength(data)).toBe(5); }); }); diff --git a/src/infra/ws.ts b/src/infra/ws.ts index 4dafd365ab69..f9b81015b693 100644 --- a/src/infra/ws.ts +++ b/src/infra/ws.ts @@ -1,18 +1,13 @@ -// Normalizes WebSocket raw payload data to strings. -import { Buffer } from "node:buffer"; +import { rawDataToString as gatewayRawDataToString } from "@openclaw/gateway-client/websocket-data"; import type WebSocket from "ws"; -// ws emits raw payloads as buffers, ArrayBuffers, or buffer fragments. +// Keep the declaration owner stable for the shipped webhook-ingress SDK export; +// WebSocket conversion itself is canonical in @openclaw/gateway-client. export function rawDataToString( data: WebSocket.RawData, encoding: BufferEncoding = "utf8", ): string { - if (Array.isArray(data)) { - return Buffer.concat(data).toString(encoding); - } - return data instanceof ArrayBuffer - ? Buffer.from(data).toString(encoding) - : data.toString(encoding); + return gatewayRawDataToString(data, encoding); } export function rawDataByteLength(data: WebSocket.RawData): number { diff --git a/src/plugins/sdk-alias.test.ts b/src/plugins/sdk-alias.test.ts index 3378ad16cd72..5aee9dc45c5d 100644 --- a/src/plugins/sdk-alias.test.ts +++ b/src/plugins/sdk-alias.test.ts @@ -1151,6 +1151,7 @@ describe("plugin sdk alias helpers", () => { const workspaceAliases = writeWorkspaceAliasFixtures(fixture.root, [ ["@openclaw/gateway-client", "gateway-client", "index"], ["@openclaw/gateway-client/timeouts", "gateway-client", "timeouts"], + ["@openclaw/gateway-client/websocket-data", "gateway-client", "websocket-data"], ["@openclaw/gateway-protocol", "gateway-protocol", "index"], ["@openclaw/gateway-protocol/schema", "gateway-protocol", "schema"], ["@openclaw/gateway-protocol/frame-guards", "gateway-protocol", "frame-guards"], diff --git a/src/plugins/sdk-alias.ts b/src/plugins/sdk-alias.ts index 148931838be4..9df3b3f9ab7e 100644 --- a/src/plugins/sdk-alias.ts +++ b/src/plugins/sdk-alias.ts @@ -489,7 +489,7 @@ const JS_STATIC_RELATIVE_DEPENDENCY_PATTERN = // Packaged installs omit workspace manifests; preserve the exact curated subpaths // instead of expanding aliases from package exports. const WORKSPACE_PACKAGE_ALIAS_SUBPATHS = [ - ["gateway-client", ["", "readiness", "timeouts"]], + ["gateway-client", ["", "readiness", "timeouts", "websocket-data"]], [ "gateway-protocol", [ diff --git a/src/worker/worker-connection-admission.ts b/src/worker/worker-connection-admission.ts index 0a4b9274a89c..b603389e8496 100644 --- a/src/worker/worker-connection-admission.ts +++ b/src/worker/worker-connection-admission.ts @@ -1,4 +1,5 @@ import { randomUUID } from "node:crypto"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { Value } from "typebox/value"; import { WebSocket, type RawData } from "ws"; import { @@ -12,7 +13,6 @@ import { } from "../../packages/gateway-protocol/src/schema/worker-admission.js"; import { WORKER_PROTOCOL_MAX_INFERENCE_PAYLOAD_BYTES } from "../../packages/gateway-protocol/src/schema/worker-inference.js"; import { PROTOCOL_VERSION } from "../../packages/gateway-protocol/src/version.js"; -import { rawDataToString } from "../infra/ws.js"; import { WorkerAdmissionError, WorkerConnectionInterruptedError, diff --git a/src/worker/worker.fault-injection.test.ts b/src/worker/worker.fault-injection.test.ts index 471351e5ab0c..a64226eed449 100644 --- a/src/worker/worker.fault-injection.test.ts +++ b/src/worker/worker.fault-injection.test.ts @@ -2,6 +2,7 @@ import fs from "node:fs/promises"; import { createServer, type Server } from "node:http"; import os from "node:os"; import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { WebSocket, WebSocketServer, type RawData } from "ws"; import { @@ -49,7 +50,6 @@ import { clearAgentRunContext, getAgentRunContext, } from "../infra/agent-run-registry.js"; -import { rawDataToString } from "../infra/ws.js"; import type { WorkerProvider, WorkerSshEndpoint } from "../plugins/types.js"; import { closeOpenClawStateDatabaseForTest, diff --git a/src/worker/worker.runtime.test.ts b/src/worker/worker.runtime.test.ts index f8eddfc3cc54..8e6f40f13880 100644 --- a/src/worker/worker.runtime.test.ts +++ b/src/worker/worker.runtime.test.ts @@ -2,6 +2,7 @@ import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { createServer, type Server } from "node:http"; import { tmpdir } from "node:os"; import path from "node:path"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { Value } from "typebox/value"; import { afterEach, describe, expect, it, vi } from "vitest"; import { WebSocket, WebSocketServer, type RawData } from "ws"; @@ -32,7 +33,6 @@ import { type WorkerInferenceTerminalOutcome, } from "../../packages/gateway-protocol/src/schema/worker-inference.js"; import { listRunningSessions } from "../agents/bash-process-registry.js"; -import { rawDataToString } from "../infra/ws.js"; import { buildWorkerConnectParams, type WorkerLaunchDescriptor } from "./launch-descriptor.js"; import { WORKER_PROVIDER_REPLAY_LOCAL_RETRY_MESSAGE } from "./transcript-message.js"; import { WorkerAdmissionDeadlineExceededError } from "./worker-connection-contract.js"; diff --git a/test/e2e/qa-lab/runtime/gateway-client-transport-defaults.e2e.test.ts b/test/e2e/qa-lab/runtime/gateway-client-transport-defaults.e2e.test.ts index 6ff442c47313..9869ceb100db 100644 --- a/test/e2e/qa-lab/runtime/gateway-client-transport-defaults.e2e.test.ts +++ b/test/e2e/qa-lab/runtime/gateway-client-transport-defaults.e2e.test.ts @@ -1,4 +1,5 @@ import { setImmediate as waitForImmediate } from "node:timers/promises"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { PROTOCOL_VERSION } from "@openclaw/gateway-protocol/version"; import { afterEach, describe, expect, it, vi } from "vitest"; import { type RawData, type WebSocket, WebSocketServer } from "ws"; @@ -19,15 +20,6 @@ type RequestFrame = { const clients: GatewayClient[] = []; let server: WebSocketServer | undefined; -function rawDataToString(data: RawData): string { - if (Array.isArray(data)) { - return Buffer.concat(data).toString("utf8"); - } - return Buffer.isBuffer(data) - ? data.toString("utf8") - : Buffer.from(new Uint8Array(data)).toString("utf8"); -} - function parseRequest(data: RawData): RequestFrame { return JSON.parse(rawDataToString(data)) as RequestFrame; } diff --git a/test/e2e/qa-lab/runtime/gateway-loopback-lan-access.ts b/test/e2e/qa-lab/runtime/gateway-loopback-lan-access.ts index b336ab9c1e1c..225612bdecb4 100644 --- a/test/e2e/qa-lab/runtime/gateway-loopback-lan-access.ts +++ b/test/e2e/qa-lab/runtime/gateway-loopback-lan-access.ts @@ -6,6 +6,7 @@ import net from "node:net"; import os from "node:os"; import path from "node:path"; import { pathToFileURL } from "node:url"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; import { WebSocket, type RawData } from "ws"; import { PROTOCOL_VERSION } from "../../../../packages/gateway-protocol/src/index.js"; import { clearConfigCache, clearRuntimeConfigSnapshot } from "../../../../src/config/config.js"; @@ -15,7 +16,6 @@ import { startGatewayServer, type GatewayServer } from "../../../../src/gateway/ import { getFreeGatewayPort } from "../../../../src/gateway/test-helpers.e2e.js"; import { GATEWAY_STARTUP_MUTATED_ENV_KEYS } from "../../../../src/gateway/test-helpers.env.js"; import { resetAgentEventsForTest } from "../../../../src/infra/agent-events.js"; -import { rawDataToString } from "../../../../src/infra/ws.js"; import { captureEnv, deleteTestEnvValue, setTestEnvValue } from "../../../../src/test-utils/env.js"; import { GATEWAY_CLIENT_MODES, diff --git a/test/scripts/gateway-ws-client.test.ts b/test/scripts/gateway-ws-client.test.ts index c296524d6656..129595213aa2 100644 --- a/test/scripts/gateway-ws-client.test.ts +++ b/test/scripts/gateway-ws-client.test.ts @@ -3,7 +3,7 @@ import { createServer, type Server } from "node:http"; import type { Duplex } from "node:stream"; import { afterEach, describe, expect, it } from "vitest"; import { WebSocket, WebSocketServer } from "ws"; -import { createGatewayWsClient, rawDataToString } from "../../scripts/dev/gateway-ws-client.js"; +import { createGatewayWsClient } from "../../scripts/dev/gateway-ws-client.js"; let server: Server | undefined; let wss: WebSocketServer | undefined; @@ -88,14 +88,6 @@ async function listenStalledUpgrade(): Promise<{ close: () => Promise; url } describe("createGatewayWsClient", () => { - it("decodes every ws raw-data shape without core source files", () => { - expect(rawDataToString(Buffer.from("buffer"))).toBe("buffer"); - expect(rawDataToString(Uint8Array.from(Buffer.from("array-buffer")).buffer)).toBe( - "array-buffer", - ); - expect(rawDataToString([Buffer.from("frag"), Buffer.from("ments")])).toBe("fragments"); - }); - it("rejects pending RPC requests when the client closes", async () => { const url = await listen(() => {}); const client = createGatewayWsClient({ url }); diff --git a/test/vitest/vitest.shared.config.ts b/test/vitest/vitest.shared.config.ts index 0e34ab7163af..fe6f7492f971 100644 --- a/test/vitest/vitest.shared.config.ts +++ b/test/vitest/vitest.shared.config.ts @@ -226,6 +226,10 @@ export const sharedVitestConfig = { find: "@openclaw/gateway-client/timeouts", replacement: path.join(repoRoot, "packages", "gateway-client", "src", "timeouts.ts"), }, + { + find: "@openclaw/gateway-client/websocket-data", + replacement: path.join(repoRoot, "packages", "gateway-client", "src", "websocket-data.ts"), + }, { find: "@openclaw/gateway-client", replacement: path.join(repoRoot, "packages", "gateway-client", "src", "index.ts"), diff --git a/tsconfig.json b/tsconfig.json index a42c3a7ee9a6..5eead2bc5603 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -70,6 +70,9 @@ "@openclaw/model-catalog-core/*": ["./packages/model-catalog-core/src/*"], "@openclaw/gateway-client": ["./packages/gateway-client/src/index.ts"], "@openclaw/gateway-client/browser": ["./packages/gateway-client/src/browser.ts"], + "@openclaw/gateway-client/websocket-data": [ + "./packages/gateway-client/src/websocket-data.ts" + ], "@openclaw/gateway-client/*": ["./packages/gateway-client/src/*"], "@openclaw/gateway-protocol": ["./packages/gateway-protocol/src/index.ts"], "@openclaw/gateway-protocol/client-info": [