From c8a99f7aabec44b2e8c0d15406ca3e35d5ba79d9 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sun, 9 Aug 2026 18:45:08 +0800 Subject: [PATCH] test(gateway): cover rolling node compatibility (#119991) Punchcard-Session: calm-lantern-timber-wa --- .../runtime/gateway-node-control-plane.yaml | 5 +- .../gateway-node-control-plane.e2e.test.ts | 837 +++++++++++++++++- 2 files changed, 824 insertions(+), 18 deletions(-) diff --git a/qa/scenarios/runtime/gateway-node-control-plane.yaml b/qa/scenarios/runtime/gateway-node-control-plane.yaml index 44db99f9c944..d4154116a912 100644 --- a/qa/scenarios/runtime/gateway-node-control-plane.yaml +++ b/qa/scenarios/runtime/gateway-node-control-plane.yaml @@ -12,12 +12,15 @@ scenario: - gateway.node-actions - gateway.node-events - gateway.remote-device-capabilities + - gateway.version-negotiation objective: Verify a real authenticated Gateway pairs and operates a remote device through the node WebSocket control plane. successCriteria: - A real child Gateway accepts an authenticated operator and pairs an isolated iOS node identity over WebSocket. - Approved reconnect declarations appear in node.list and node.describe with the effective capabilities, commands, permissions, name, platform, and connected state. - Exact camera.list and location.get requests and results cross the same node.invoke and node.invoke.result socket boundary. - A node.presence.alive event is persisted and becomes visible through both node.list and node.describe. + - A synthetic overlap fixture records the default node-role and node-mode v3-through-v4 envelope and rejects the current-only operator envelope; shipped-v3 peer behavior is outside this scenario. + - The same paired node identity reconnects across v3-only and v4 envelopes, with plugin features hidden and restored according to the advertised envelope. docsRefs: - docs/gateway/protocol.md - docs/nodes/index.md @@ -29,4 +32,4 @@ scenario: execution: kind: vitest path: test/e2e/qa-lab/runtime/gateway-node-control-plane.e2e.test.ts - summary: Start a real child Gateway, pair an iOS node WebSocket, inspect inventory and declarations, invoke camera and location commands, and persist node presence. + summary: Capture the default node v3-through-v4 connect envelope with a synthetic overlap fixture, pair a current iOS identity to a real child Gateway, verify v3-only filtering and same-identity reconnect behavior, invoke built-in and plugin commands, and persist node presence. diff --git a/test/e2e/qa-lab/runtime/gateway-node-control-plane.e2e.test.ts b/test/e2e/qa-lab/runtime/gateway-node-control-plane.e2e.test.ts index a154551c6f1a..bbfe8e0ce6fd 100644 --- a/test/e2e/qa-lab/runtime/gateway-node-control-plane.e2e.test.ts +++ b/test/e2e/qa-lab/runtime/gateway-node-control-plane.e2e.test.ts @@ -1,18 +1,32 @@ // Proves the Gateway node control plane across real authenticated WebSockets. import { randomUUID } from "node:crypto"; import { existsSync } from "node:fs"; +import fs from "node:fs/promises"; +import os from "node:os"; import path from "node:path"; +import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { GatewayClient } from "openclaw/plugin-sdk/gateway-runtime"; -import { describe, expect, it, vi } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { type RawData, WebSocketServer } from "ws"; import { startQaGatewayChild } from "../../../../extensions/qa-lab/api.js"; import { GATEWAY_CLIENT_MODES, GATEWAY_CLIENT_NAMES, } from "../../../../packages/gateway-protocol/src/client-info.js"; +import { ConnectErrorDetailCodes } from "../../../../packages/gateway-protocol/src/connect-error-details.js"; +import { + ErrorCodes, + MIN_CLIENT_PROTOCOL_VERSION, + MIN_NODE_PROTOCOL_VERSION, + PROTOCOL_VERSION, + type HelloOk, +} from "../../../../packages/gateway-protocol/src/index.js"; +import type { OpenClawConfig } from "../../../../src/config/types.openclaw.js"; import { loadOrCreateDeviceIdentity, type DeviceIdentity, } from "../../../../src/infra/device-identity.js"; +import { useAutoCleanupTempDirTracker } from "../../../helpers/temp-dir.js"; const TEST_TIMEOUT_MS = 180_000; const REQUEST_TIMEOUT_MS = 20_000; @@ -24,8 +38,13 @@ const NODE_PERMISSIONS = { camera: true, location: true, }; +const FIXTURE_PLUGIN_ID = "qa-gateway-node-rolling-compat"; +const FIXTURE_CAPABILITY = "qa-rolling-surface"; +const FIXTURE_COMMAND = "qa.rolling.echo"; +const FIXTURE_ROUTE = "/qa-rolling-surface"; type GatewayHandle = Awaited>; +type GatewayConnection = Pick; type NodeRead = { nodeId: string; displayName?: string; @@ -52,8 +71,23 @@ type InvocationRecord = { command: string; params: unknown; }; +type ConnectEnvelope = { + role: string; + mode: string; + clientName: string; + platform: string; + deviceFamily?: string; + minProtocol: number; + maxProtocol: number; +}; +type StoppableClient = Pick; +type StoppableFixture = { + stop: () => Promise; +}; describe("Gateway node control plane", () => { + const tempDirs = useAutoCleanupTempDirTracker(afterEach); + it( "pairs, inventories, invokes, and records presence for one remote device", { timeout: TEST_TIMEOUT_MS }, @@ -93,6 +127,7 @@ describe("Gateway node control plane", () => { }); const invocations: InvocationRecord[] = []; const handlerErrors: Error[] = []; + const invocationResponses: Promise[] = []; let operator: GatewayClient | undefined; let node: GatewayClient | undefined; @@ -106,9 +141,12 @@ describe("Gateway node control plane", () => { if (event.event !== "node.invoke.request") { return; } - void respondToInvocation(node, event.payload, invocations).catch((error: unknown) => { - handlerErrors.push(error instanceof Error ? error : new Error(String(error))); - }); + const response = respondToInvocation(node, event.payload, invocations).catch( + (error: unknown) => { + handlerErrors.push(error instanceof Error ? error : new Error(String(error))); + }, + ); + invocationResponses.push(response); }, }); @@ -210,6 +248,7 @@ describe("Gateway node control plane", () => { params: locationParams, }, ]); + await Promise.all(invocationResponses); expect(handlerErrors).toEqual([]); const aliveSentAtMs = Date.now(); @@ -260,6 +299,7 @@ describe("Gateway node control plane", () => { { timeout: REQUEST_TIMEOUT_MS, interval: 100 }, ); } finally { + await Promise.all(invocationResponses); await Promise.allSettled([ ...(node ? [node.stopAndWait({ timeoutMs: 1_000 })] : []), ...(operator ? [operator.stopAndWait({ timeoutMs: 1_000 })] : []), @@ -270,9 +310,430 @@ describe("Gateway node control plane", () => { } }, ); + + it( + "captures the default node protocol envelope with a synthetic v3 overlap fixture", + { timeout: REQUEST_TIMEOUT_MS * 2 }, + async () => { + expect(MIN_CLIENT_PROTOCOL_VERSION).toBe(PROTOCOL_VERSION); + expect(MIN_NODE_PROTOCOL_VERSION).toBeLessThan(PROTOCOL_VERSION); + + const fixture = await startProtocolEnvelopeFixture(); + let node: GatewayClient | undefined; + let fixtureHello: HelloOk | undefined; + try { + node = await connectClient({ + gateway: fixture.gateway, + role: "node", + clientName: GATEWAY_CLIENT_NAMES.IOS_APP, + clientDisplayName: NODE_DISPLAY_NAME, + mode: GATEWAY_CLIENT_MODES.NODE, + platform: "ios", + deviceFamily: "iPhone", + scopes: [], + deviceIdentity: null, + onHelloOk: (value) => { + fixtureHello = value; + }, + }); + + // hello.protocol reports the fixture Gateway's current server protocol. + expect(fixtureHello?.protocol).toBe(MIN_NODE_PROTOCOL_VERSION); + expect(fixture.connectFrames).toMatchObject([ + { + role: "node", + mode: GATEWAY_CLIENT_MODES.NODE, + minProtocol: MIN_NODE_PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + }, + ]); + + await expect(connectOperator(fixture.gateway)).rejects.toThrow(/protocol mismatch/i); + expect(fixture.connectFrames).toMatchObject([ + { + role: "node", + mode: GATEWAY_CLIENT_MODES.NODE, + minProtocol: MIN_NODE_PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + }, + { + role: "operator", + mode: GATEWAY_CLIENT_MODES.BACKEND, + minProtocol: MIN_CLIENT_PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + }, + ]); + } finally { + await stopProtocolEnvelopeFixture(node, fixture); + } + }, + ); + + it( + "keeps one node-host client converged across Gateway upgrades and rollbacks", + { timeout: REQUEST_TIMEOUT_MS * 4 }, + async () => { + const identity = loadOrCreateDeviceIdentity({ + path: path.join(tempDirs.make("openclaw-node-host-negotiation-"), "node.sqlite"), + }); + const fixture = await startProtocolEnvelopeFixture(); + const helloProtocols: number[] = []; + let node: GatewayClient | undefined; + + try { + // runNodeHost owns one NODE_HOST GatewayClient. Keep this exact instance + // across server transitions so production negotiation owns every retry. + node = await connectClient({ + gateway: fixture.gateway, + role: "node", + clientName: GATEWAY_CLIENT_NAMES.NODE_HOST, + clientDisplayName: "QA node host", + mode: GATEWAY_CLIENT_MODES.NODE, + platform: "macos", + deviceFamily: "Mac", + scopes: [], + deviceIdentity: identity, + onHelloOk: (hello) => { + helloProtocols.push(hello.protocol); + }, + }); + + expect(helloProtocols).toEqual([MIN_NODE_PROTOCOL_VERSION]); + expect(fixture.connectFrames).toMatchObject([ + { + clientName: GATEWAY_CLIENT_NAMES.NODE_HOST, + platform: "macos", + deviceFamily: "Mac", + minProtocol: PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + }, + { + clientName: GATEWAY_CLIENT_NAMES.NODE_HOST, + platform: "darwin", + minProtocol: MIN_NODE_PROTOCOL_VERSION, + maxProtocol: MIN_NODE_PROTOCOL_VERSION, + }, + ]); + expect(fixture.connectFrames[1]).not.toHaveProperty("deviceFamily"); + + fixture.setProtocol(PROTOCOL_VERSION); + fixture.disconnectClients("Gateway upgraded to protocol v4"); + await vi.waitFor( + () => expect(helloProtocols).toEqual([MIN_NODE_PROTOCOL_VERSION, PROTOCOL_VERSION]), + { + timeout: REQUEST_TIMEOUT_MS, + interval: 100, + }, + ); + expect(fixture.connectFrames.slice(-2)).toMatchObject([ + { + platform: "darwin", + minProtocol: MIN_NODE_PROTOCOL_VERSION, + maxProtocol: MIN_NODE_PROTOCOL_VERSION, + }, + { + platform: "macos", + deviceFamily: "Mac", + minProtocol: PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + }, + ]); + + fixture.setProtocol(MIN_NODE_PROTOCOL_VERSION); + fixture.disconnectClients("Gateway rolled back to protocol v3"); + await vi.waitFor( + () => + expect(helloProtocols).toEqual([ + MIN_NODE_PROTOCOL_VERSION, + PROTOCOL_VERSION, + MIN_NODE_PROTOCOL_VERSION, + ]), + { + timeout: REQUEST_TIMEOUT_MS, + interval: 100, + }, + ); + expect(fixture.connectFrames.slice(-2)).toMatchObject([ + { + platform: "macos", + deviceFamily: "Mac", + minProtocol: PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + }, + { + platform: "darwin", + minProtocol: MIN_NODE_PROTOCOL_VERSION, + maxProtocol: MIN_NODE_PROTOCOL_VERSION, + }, + ]); + } finally { + await stopProtocolEnvelopeFixture(node, fixture); + } + }, + ); + + it("stops the protocol fixture when client cleanup rejects", async () => { + const clientError = new Error("client cleanup failed"); + const stopFixture = vi.fn(async () => {}); + + await expect( + stopProtocolEnvelopeFixture( + { + stopAndWait: vi.fn(async () => { + throw clientError; + }), + }, + { stop: stopFixture }, + ), + ).rejects.toBe(clientError); + expect(stopFixture).toHaveBeenCalledOnce(); + }); + + it( + "reconnects the same paired identity across v3-only and v4 envelopes", + { timeout: TEST_TIMEOUT_MS }, + async () => { + expect(MIN_NODE_PROTOCOL_VERSION).toBe(3); + expect(PROTOCOL_VERSION).toBe(4); + + const fixture = await createFixturePlugin(); + let gateway: GatewayHandle | undefined; + let operator: GatewayClient | undefined; + let node: GatewayClient | undefined; + let proofError: unknown; + const cleanupErrors: unknown[] = []; + const invocations: InvocationRecord[] = []; + const handlerErrors: Error[] = []; + const invocationResponses: Promise[] = []; + + try { + gateway = await startQaGatewayChild({ + repoRoot: process.cwd(), + command: { + executablePath: process.execPath, + argsPrefix: ["--import", "tsx", "src/entry.ts"], + cwd: process.cwd(), + usePackagedPlugins: true, + }, + transportBaseUrl: "http://127.0.0.1", + controlUiEnabled: false, + runtimeEnvPatch: { + OPENCLAW_DISABLE_BUNDLED_PLUGINS: "1", + OPENCLAW_SKIP_CHANNELS: "1", + OPENCLAW_SKIP_PROVIDERS: "1", + OPENCLAW_TEST_MINIMAL_GATEWAY: "1", + }, + mutateConfig: (cfg) => { + // Bundled plugins are disabled for this focused proof, so their + // configured slots cannot remain while the fixture plugin is merged. + const { slots: _bundledSlots, ...plugins } = cfg.plugins ?? {}; + return withFixturePlugin( + { + ...cfg, + plugins, + gateway: { + ...cfg.gateway, + nodes: { + ...cfg.gateway?.nodes, + commands: { allow: ["camera.list", FIXTURE_COMMAND] }, + }, + }, + }, + fixture.pluginDir, + ); + }, + }); + const identity = loadOrCreateDeviceIdentity({ + path: path.join(gateway.tempRoot, "rolling-compat-node.sqlite"), + }); + const declaredCaps = ["camera", FIXTURE_CAPABILITY]; + const declaredCommands = ["camera.list", FIXTURE_COMMAND]; + const onEvent = (event: { event: string; payload?: unknown }) => { + if (event.event !== "node.invoke.request") { + return; + } + const response = respondToInvocation(node, event.payload, invocations).catch( + (error: unknown) => { + handlerErrors.push(error instanceof Error ? error : new Error(String(error))); + }, + ); + invocationResponses.push(response); + }; + + operator = await connectOperator(gateway); + let constrainedHello: HelloOk | undefined; + node = await connectPairedNode({ + gateway, + identity, + operator, + caps: declaredCaps, + commands: declaredCommands, + minProtocol: MIN_NODE_PROTOCOL_VERSION, + maxProtocol: MIN_NODE_PROTOCOL_VERSION, + onEvent, + onHelloOk: (hello) => { + constrainedHello = hello; + }, + }); + + // The client envelope is constrained to v3; hello reports the v4 server protocol. + expect(constrainedHello?.protocol).toBe(PROTOCOL_VERSION); + expect(constrainedHello?.pluginSurfaceUrls).toBeUndefined(); + const legacyListed = await waitForApprovedNode(operator, identity.deviceId, gateway.logs); + expect(legacyListed.caps).toEqual(["camera"]); + expect(legacyListed.commands).toEqual(["camera.list"]); + + const cameraParams = { includeUnavailable: false }; + const cameraResult = await invokeNodeCommand({ + operator, + nodeId: identity.deviceId, + command: "camera.list", + params: cameraParams, + }); + expect(cameraResult).toMatchObject({ + ok: true, + nodeId: identity.deviceId, + command: "camera.list", + payload: { + cameras: [{ id: "back-wide", position: "back" }], + received: cameraParams, + }, + }); + + await Promise.all(invocationResponses); + expect(handlerErrors).toEqual([]); + await node.stopAndWait({ timeoutMs: 1_000 }); + node = undefined; + await expectNoPendingPairing(operator, identity.deviceId); + + let currentHello: HelloOk | undefined; + node = await connectClient({ + gateway, + role: "node", + clientName: GATEWAY_CLIENT_NAMES.IOS_APP, + clientDisplayName: NODE_DISPLAY_NAME, + mode: GATEWAY_CLIENT_MODES.NODE, + platform: "ios", + deviceFamily: "iPhone", + scopes: [], + caps: declaredCaps, + commands: declaredCommands, + permissions: NODE_PERMISSIONS, + deviceIdentity: identity, + minProtocol: PROTOCOL_VERSION, + maxProtocol: PROTOCOL_VERSION, + onEvent, + onHelloOk: (hello) => { + currentHello = hello; + }, + }); + + expect(currentHello?.protocol).toBe(PROTOCOL_VERSION); + const fixtureSurfaceUrl = currentHello?.pluginSurfaceUrls?.[FIXTURE_CAPABILITY]; + if (!fixtureSurfaceUrl) { + throw new Error("v4 hello omitted the fixture plugin surface URL"); + } + expect(new URL(fixtureSurfaceUrl).pathname).toMatch(/^\/__openclaw__\/cap\//); + const fixtureSurfaceResponse = await fetch(`${fixtureSurfaceUrl}${FIXTURE_ROUTE}`, { + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }); + expect(fixtureSurfaceResponse.status).toBe(204); + await expectNoPendingPairing(operator, identity.deviceId); + const currentListed = await waitForConnectedApprovedNode( + operator, + identity.deviceId, + gateway.logs, + ); + expect(currentListed.caps?.toSorted()).toEqual(declaredCaps.toSorted()); + expect(currentListed.commands?.toSorted()).toEqual(declaredCommands.toSorted()); + + const pluginParams = { message: "rolling-compatible" }; + const pluginResult = await invokeNodeCommand({ + operator, + nodeId: identity.deviceId, + command: FIXTURE_COMMAND, + params: pluginParams, + }); + expect(pluginResult).toMatchObject({ + ok: true, + nodeId: identity.deviceId, + command: FIXTURE_COMMAND, + payload: { + echoed: pluginParams, + }, + }); + expect(invocations).toMatchObject([ + { + nodeId: identity.deviceId, + command: "camera.list", + params: cameraParams, + }, + { + nodeId: identity.deviceId, + command: FIXTURE_COMMAND, + params: pluginParams, + }, + ]); + await Promise.all(invocationResponses); + expect(handlerErrors).toEqual([]); + } catch (error) { + proofError = error; + } finally { + await Promise.all(invocationResponses); + const clientCleanup = await Promise.allSettled([ + ...(node ? [node.stopAndWait({ timeoutMs: 1_000 })] : []), + ...(operator ? [operator.stopAndWait({ timeoutMs: 1_000 })] : []), + ]); + for (const result of clientCleanup) { + if (result.status === "rejected") { + cleanupErrors.push(result.reason); + } + } + if (gateway) { + const tempRoot = gateway.tempRoot; + try { + await gateway.stop(); + expect(existsSync(tempRoot)).toBe(false); + } catch (error) { + cleanupErrors.push(error); + } + } + try { + await fixture.cleanup(); + expect(existsSync(fixture.root)).toBe(false); + } catch (error) { + cleanupErrors.push(error); + } + } + const failures = proofError === undefined ? cleanupErrors : [proofError, ...cleanupErrors]; + if (failures.length === 1) { + throw failures[0]; + } + if (failures.length > 1) { + throw new AggregateError(failures, "gateway node rolling compatibility proof failed"); + } + }, + ); }); -async function connectOperator(gateway: GatewayHandle): Promise { +async function stopProtocolEnvelopeFixture( + node: StoppableClient | undefined, + fixture: StoppableFixture, +): Promise { + const results = await Promise.allSettled([ + ...(node ? [node.stopAndWait({ timeoutMs: 1_000 })] : []), + fixture.stop(), + ]); + const errors = results.flatMap((result) => (result.status === "rejected" ? [result.reason] : [])); + if (errors.length === 1) { + throw errors[0]; + } + if (errors.length > 1) { + throw new AggregateError(errors, "protocol envelope fixture cleanup failed"); + } +} + +async function connectOperator(gateway: GatewayConnection): Promise { return await connectClient({ gateway, role: "operator", @@ -285,10 +746,15 @@ async function connectOperator(gateway: GatewayHandle): Promise { } async function connectPairedNode(params: { - gateway: GatewayHandle; + gateway: GatewayConnection; identity: DeviceIdentity; operator: GatewayClient; + caps?: string[]; + commands?: string[]; + minProtocol?: number; + maxProtocol?: number; onEvent: (event: { event: string; payload?: unknown }) => void; + onHelloOk?: (hello: HelloOk) => void; }): Promise { const connect = () => connectClient({ @@ -300,11 +766,14 @@ async function connectPairedNode(params: { platform: "ios", deviceFamily: "iPhone", scopes: [], - caps: NODE_CAPS, - commands: NODE_COMMANDS, + caps: params.caps ?? NODE_CAPS, + commands: params.commands ?? NODE_COMMANDS, permissions: NODE_PERMISSIONS, deviceIdentity: params.identity, + minProtocol: params.minProtocol, + maxProtocol: params.maxProtocol, onEvent: params.onEvent, + onHelloOk: params.onHelloOk, }); try { return await connect(); @@ -318,9 +787,12 @@ async function connectPairedNode(params: { } async function connectClient(params: { - gateway: GatewayHandle; + gateway: GatewayConnection; role: "operator" | "node"; - clientName: typeof GATEWAY_CLIENT_NAMES.GATEWAY_CLIENT | typeof GATEWAY_CLIENT_NAMES.IOS_APP; + clientName: + | typeof GATEWAY_CLIENT_NAMES.GATEWAY_CLIENT + | typeof GATEWAY_CLIENT_NAMES.IOS_APP + | typeof GATEWAY_CLIENT_NAMES.NODE_HOST; clientDisplayName: string; mode: typeof GATEWAY_CLIENT_MODES.BACKEND | typeof GATEWAY_CLIENT_MODES.NODE; scopes: string[]; @@ -330,7 +802,10 @@ async function connectClient(params: { commands?: string[]; permissions?: Record; deviceIdentity: DeviceIdentity | null; + minProtocol?: number; + maxProtocol?: number; onEvent?: (event: { event: string; payload?: unknown }) => void; + onHelloOk?: (hello: HelloOk) => void; }): Promise { return await new Promise((resolve, reject) => { let settled = false; @@ -363,11 +838,27 @@ async function connectClient(params: { commands: params.commands, permissions: params.permissions, deviceIdentity: params.deviceIdentity, + minProtocol: params.minProtocol, + maxProtocol: params.maxProtocol, requestTimeoutMs: REQUEST_TIMEOUT_MS, onEvent: params.onEvent, - onHelloOk: () => finish(), + onHelloOk: (hello) => { + params.onHelloOk?.(hello); + finish(); + }, onConnectError: (error) => finish(error), - onClose: (code, reason) => finish(new Error(`Gateway closed (${code}): ${reason}`)), + onClose: (code, reason) => { + // NODE_HOST uses this exact close to switch protocol envelopes. Keep the + // same client alive until its negotiated hello reaches the caller. + if ( + params.clientName === GATEWAY_CLIENT_NAMES.NODE_HOST && + code === 1008 && + reason === "connect retry" + ) { + return; + } + finish(new Error(`Gateway closed (${code}): ${reason}`)); + }, }); const timeout = setTimeout( () => finish(new Error(`Gateway client connection timed out:\n${params.gateway.logs()}`)), @@ -451,11 +942,32 @@ async function waitForApprovedNode( return approved; } +async function waitForConnectedApprovedNode( + operator: GatewayClient, + nodeId: string, + logs: () => string, +): Promise { + let approved: NodeRead | undefined; + await vi.waitFor( + async () => { + approved = await readNode(operator, nodeId); + expect(approved, logs()).toMatchObject({ + nodeId, + approvalState: "approved", + connected: true, + paired: true, + }); + }, + { timeout: REQUEST_TIMEOUT_MS, interval: 100 }, + ); + if (!approved) { + throw new Error(`approved node never became visible:\n${logs()}`); + } + return approved; +} + async function approvePendingNodeSurface(operator: GatewayClient, nodeId: string): Promise { - const nodes = await operator.request<{ - pending?: Array<{ requestId?: string; nodeId?: string }>; - }>("node.pair.list", {}, { timeoutMs: REQUEST_TIMEOUT_MS }); - for (const pending of nodes.pending ?? []) { + for (const pending of await readPendingNodePairings(operator)) { if (pending.nodeId === nodeId && pending.requestId) { await operator.request( "node.pair.approve", @@ -466,6 +978,28 @@ async function approvePendingNodeSurface(operator: GatewayClient, nodeId: string } } +async function readPendingNodePairings( + operator: GatewayClient, +): Promise> { + const nodes = await operator.request<{ + pending?: Array<{ requestId?: string; nodeId?: string }>; + }>("node.pair.list", {}, { timeoutMs: REQUEST_TIMEOUT_MS }); + return nodes.pending ?? []; +} + +async function expectNoPendingPairing(operator: GatewayClient, nodeId: string): Promise { + const [devices, nodes] = await Promise.all([ + operator.request<{ + pending?: Array<{ deviceId?: string }>; + }>("device.pair.list", {}, { timeoutMs: REQUEST_TIMEOUT_MS }), + readPendingNodePairings(operator), + ]); + const devicePending = devices.pending?.some((entry) => entry.deviceId === nodeId) ?? false; + const nodePending = nodes.some((entry) => entry.nodeId === nodeId); + expect(devicePending).toBe(false); + expect(nodePending).toBe(false); +} + async function readNode(operator: GatewayClient, nodeId: string): Promise { const result = await operator.request<{ nodes?: NodeRead[] }>( "node.list", @@ -475,6 +1009,30 @@ async function readNode(operator: GatewayClient, nodeId: string): Promise entry.nodeId === nodeId); } +async function invokeNodeCommand(params: { + operator: GatewayClient; + nodeId: string; + command: string; + params: unknown; +}): Promise<{ + ok: boolean; + nodeId: string; + command: string; + payload: unknown; +}> { + return await params.operator.request( + "node.invoke", + { + nodeId: params.nodeId, + command: params.command, + params: params.params, + timeoutMs: REQUEST_TIMEOUT_MS, + idempotencyKey: randomUUID(), + }, + { timeoutMs: REQUEST_TIMEOUT_MS }, + ); +} + async function respondToInvocation( node: GatewayClient | undefined, payload: unknown, @@ -503,7 +1061,11 @@ async function respondToInvocation( longitude: -122.0312, received: params, } - : undefined; + : frame.command === FIXTURE_COMMAND + ? { + echoed: params, + } + : undefined; if (!response) { throw new Error(`unexpected node command: ${frame.command}`); } @@ -518,3 +1080,244 @@ async function respondToInvocation( { timeoutMs: REQUEST_TIMEOUT_MS }, ); } + +async function createFixturePlugin(): Promise<{ + root: string; + pluginDir: string; + cleanup: () => Promise; +}> { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-gateway-node-rolling-")); + const pluginDir = path.join(root, FIXTURE_PLUGIN_ID); + try { + await fs.mkdir(pluginDir, { recursive: true }); + await fs.writeFile( + path.join(pluginDir, "openclaw.plugin.json"), + `${JSON.stringify( + { + id: FIXTURE_PLUGIN_ID, + activation: { onStartup: true }, + configSchema: { type: "object", additionalProperties: false, properties: {} }, + }, + null, + 2, + )}\n`, + "utf8", + ); + await fs.writeFile( + path.join(pluginDir, "index.js"), + `module.exports = { + id: ${JSON.stringify(FIXTURE_PLUGIN_ID)}, + register(api) { + api.registerHttpRoute({ + path: ${JSON.stringify(FIXTURE_ROUTE)}, + auth: "plugin", + nodeCapability: { surface: ${JSON.stringify(FIXTURE_CAPABILITY)} }, + handler(_req, res) { + res.statusCode = 204; + res.end(); + return true; + }, + }); + api.registerNodeInvokePolicy({ + commands: [${JSON.stringify(FIXTURE_COMMAND)}], + defaultPlatforms: ["ios"], + handle: async (ctx) => await ctx.invokeNode(), + }); + }, +};\n`, + "utf8", + ); + return { + root, + pluginDir, + cleanup: () => fs.rm(root, { force: true, recursive: true }), + }; + } catch (error) { + try { + await fs.rm(root, { force: true, recursive: true }); + } catch (cleanupError) { + const failure = new AggregateError( + [error, cleanupError], + "fixture plugin setup and cleanup failed", + ); + failure.cause = error; + throw failure; + } + throw error; + } +} + +function withFixturePlugin(config: OpenClawConfig, pluginDir: string): OpenClawConfig { + return { + ...config, + plugins: { + ...config.plugins, + enabled: true, + allow: [...new Set([...(config.plugins?.allow ?? []), FIXTURE_PLUGIN_ID])], + load: { + ...config.plugins?.load, + paths: [...new Set([...(config.plugins?.load?.paths ?? []), pluginDir])], + }, + entries: { + ...config.plugins?.entries, + [FIXTURE_PLUGIN_ID]: { enabled: true }, + }, + }, + }; +} + +function parseConnectEnvelope( + data: RawData, +): { id: string; envelope: ConnectEnvelope } | undefined { + try { + const text = Array.isArray(data) + ? Buffer.concat(data.map((chunk) => Buffer.from(chunk))).toString("utf8") + : Buffer.isBuffer(data) + ? data.toString("utf8") + : Buffer.from(data).toString("utf8"); + const frame = JSON.parse(text); + if (!isRecord(frame) || frame.method !== "connect" || typeof frame.id !== "string") { + return undefined; + } + const params = frame.params; + if (!isRecord(params) || !isRecord(params.client)) { + return undefined; + } + if ( + typeof params.role !== "string" || + typeof params.client.id !== "string" || + typeof params.client.mode !== "string" || + typeof params.client.platform !== "string" || + typeof params.minProtocol !== "number" || + typeof params.maxProtocol !== "number" + ) { + return undefined; + } + return { + id: frame.id, + envelope: { + role: params.role, + mode: params.client.mode, + clientName: params.client.id, + platform: params.client.platform, + ...(typeof params.client.deviceFamily === "string" + ? { deviceFamily: params.client.deviceFamily } + : {}), + minProtocol: params.minProtocol, + maxProtocol: params.maxProtocol, + }, + }; + } catch { + return undefined; + } +} + +async function startProtocolEnvelopeFixture(): Promise<{ + gateway: GatewayConnection; + connectFrames: ConnectEnvelope[]; + setProtocol: (protocol: number) => void; + disconnectClients: (reason: string) => void; + stop: () => Promise; +}> { + const connectFrames: ConnectEnvelope[] = []; + let protocol: number = MIN_NODE_PROTOCOL_VERSION; + let challengeSequence = 0; + const server = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + server.on("connection", (socket) => { + socket.send( + JSON.stringify({ + type: "event", + event: "connect.challenge", + seq: 1, + payload: { nonce: `qa-protocol-envelope-${++challengeSequence}`, ts: Date.now() }, + }), + ); + socket.on("message", (data) => { + const connect = parseConnectEnvelope(data); + if (!connect) { + return; + } + connectFrames.push(connect.envelope); + const supportsProtocol = + connect.envelope.minProtocol <= protocol && connect.envelope.maxProtocol >= protocol; + if (!supportsProtocol) { + socket.send( + JSON.stringify({ + type: "res", + id: connect.id, + ok: false, + error: { + code: ErrorCodes.INVALID_REQUEST, + message: "protocol mismatch", + details: { + code: ConnectErrorDetailCodes.PROTOCOL_MISMATCH, + clientMinProtocol: connect.envelope.minProtocol, + clientMaxProtocol: connect.envelope.maxProtocol, + expectedProtocol: protocol, + }, + }, + }), + ); + return; + } + socket.send( + JSON.stringify({ + type: "res", + id: connect.id, + ok: true, + payload: { + type: "hello-ok", + protocol, + server: { version: `qa-protocol-${protocol}-fixture`, connId: randomUUID() }, + features: { methods: [], events: [] }, + snapshot: { + presence: [], + health: {}, + stateVersion: { presence: 1, health: 1 }, + uptimeMs: 1, + }, + auth: { role: connect.envelope.role, scopes: [] }, + policy: { + maxPayload: 512 * 1024, + maxBufferedBytes: 1024 * 1024, + tickIntervalMs: 30_000, + }, + }, + }), + ); + }); + }); + await new Promise((resolve, reject) => { + server.once("listening", resolve); + server.once("error", reject); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("v3 Gateway fixture did not bind a TCP port"); + } + return { + connectFrames, + gateway: { + wsUrl: `ws://127.0.0.1:${address.port}`, + token: "qa-v3-gateway-token", + runtimeEnv: {}, + logs: () => JSON.stringify(connectFrames), + }, + setProtocol(nextProtocol) { + protocol = nextProtocol; + }, + disconnectClients(reason) { + for (const client of server.clients) { + client.close(1012, reason); + } + }, + async stop() { + for (const client of server.clients) { + client.terminate(); + } + await new Promise((resolve) => { + server.close(() => resolve()); + }); + }, + }; +}