mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-12 21:53:00 -06:00
refactor(ws): consolidate raw WebSocket payload decoding (#121268)
* refactor(ws): consolidate raw data conversion * fix(scripts): keep gateway client source-loadable * refactor(ws): share plugin frame decoding
This commit is contained in:
committed by
GitHub
parent
1c9649afc4
commit
bd6c6aaef2
@@ -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";
|
||||
|
||||
|
||||
@@ -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(() =>
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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());
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
@@ -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) {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
+3
-18
@@ -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)));
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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:*",
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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";
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
+4
-9
@@ -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 {
|
||||
|
||||
@@ -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"],
|
||||
|
||||
@@ -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",
|
||||
[
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<void>; 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 });
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -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": [
|
||||
|
||||
Reference in New Issue
Block a user