mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-12 21:53:00 -06:00
fix(e2e): verify suspend control over WebSocket (#122888)
This commit is contained in:
@@ -1,7 +1,6 @@
|
||||
// WebSocket client helpers for gateway network E2E scenarios.
|
||||
import assert from "node:assert/strict";
|
||||
import { readFile, writeFile } from "node:fs/promises";
|
||||
import { request as httpRequest } from "node:http";
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { WebSocket } from "ws";
|
||||
import { isRecord } from "../../../lib/record-shared.mjs";
|
||||
@@ -95,6 +94,7 @@ function hasGatewayHealthSummaryPayload(
|
||||
}
|
||||
const { payload } = response;
|
||||
return (
|
||||
response.ok === true &&
|
||||
payload.ok === true &&
|
||||
typeof payload.ts === "number" &&
|
||||
typeof payload.durationMs === "number" &&
|
||||
@@ -158,47 +158,6 @@ async function readProbe(
|
||||
return await readJson<GatewayProbeResponse["body"]>(response, pathname, signal);
|
||||
}
|
||||
|
||||
async function requestUpgradeRejection(
|
||||
url: string,
|
||||
timeoutMs: number,
|
||||
): Promise<{ body: string; status: number | undefined }> {
|
||||
const target = new URL(httpUrl(url));
|
||||
return await new Promise<{ body: string; status: number | undefined }>((resolve, reject) => {
|
||||
const request = httpRequest(
|
||||
{
|
||||
hostname: target.hostname,
|
||||
port: target.port,
|
||||
path: target.pathname,
|
||||
headers: {
|
||||
Connection: "Upgrade",
|
||||
Upgrade: "websocket",
|
||||
"Sec-WebSocket-Key": Buffer.from("gateway-net-e2e!").toString("base64"),
|
||||
"Sec-WebSocket-Version": "13",
|
||||
},
|
||||
},
|
||||
(response) => {
|
||||
const chunks: Buffer[] = [];
|
||||
response.on("data", (chunk) => chunks.push(Buffer.from(chunk)));
|
||||
response.on("end", () => {
|
||||
resolve({
|
||||
status: response.statusCode,
|
||||
body: Buffer.concat(chunks).toString("utf8"),
|
||||
});
|
||||
});
|
||||
},
|
||||
);
|
||||
request.on("upgrade", (response, socket) => {
|
||||
socket.destroy();
|
||||
reject(new Error(`expected rejected websocket upgrade, received ${response.statusCode}`));
|
||||
});
|
||||
request.on("error", reject);
|
||||
request.setTimeout(timeoutMs, () => {
|
||||
request.destroy(new Error("websocket upgrade rejection timeout"));
|
||||
});
|
||||
request.end();
|
||||
});
|
||||
}
|
||||
|
||||
function emitPhase(phase: string, startedAt: number) {
|
||||
console.log(
|
||||
JSON.stringify({
|
||||
@@ -321,6 +280,80 @@ function assertAdminSuccess(response: GatewayAdminResponse, message: string) {
|
||||
return assertRpcSuccess(response.body, message);
|
||||
}
|
||||
|
||||
export async function verifyPreparedSuspensionSocket(
|
||||
options: GatewayClientOptions & { deadline: number; suspensionId: string },
|
||||
deps: Pick<GatewayNetworkDeps, "onceFrame" | "openSocket" | "protocolVersion"> = {},
|
||||
) {
|
||||
const { deadline, suspensionId, token, url } = options;
|
||||
const onceFrameImpl = deps.onceFrame ?? onceFrame;
|
||||
const ws = await (deps.openSocket ?? openSocket)(url, remainingDeadlineMs(deadline));
|
||||
try {
|
||||
let requestIndex = 0;
|
||||
const request = async (method: string, params: Record<string, unknown> = {}) => {
|
||||
const id = `s${++requestIndex}`;
|
||||
ws.send(JSON.stringify({ type: "req", id, method, params }));
|
||||
return (await onceFrameImpl(
|
||||
ws,
|
||||
(frame) => frame?.type === "res" && frame?.id === id,
|
||||
remainingDeadlineMs(deadline),
|
||||
)) as GatewayFrame;
|
||||
};
|
||||
const protocolVersion = deps.protocolVersion ?? (await readProtocolVersion());
|
||||
assertRpcSuccess(
|
||||
await request("connect", {
|
||||
minProtocol: protocolVersion,
|
||||
maxProtocol: protocolVersion,
|
||||
client: {
|
||||
id: "cli",
|
||||
displayName: "docker-net-e2e",
|
||||
version: "dev",
|
||||
platform: process.platform,
|
||||
mode: "cli",
|
||||
},
|
||||
caps: [],
|
||||
auth: { token },
|
||||
role: "operator",
|
||||
scopes: ["operator.admin"],
|
||||
}),
|
||||
"prepared suspension connect",
|
||||
);
|
||||
const initialStatus = assertRpcSuccess(
|
||||
await request("gateway.suspend.status", { suspensionId }),
|
||||
"prepared suspension status",
|
||||
);
|
||||
assert.equal(initialStatus?.status, "ready", "prepared suspension must remain ready");
|
||||
assertGatewaySuspendingError(await request("health"));
|
||||
const wrongResume = await request("gateway.suspend.resume", {
|
||||
suspensionId: `${suspensionId}-wrong`,
|
||||
});
|
||||
assert.equal(wrongResume.ok, false, "wrong suspension id must fail");
|
||||
assert.equal(wrongResume.error?.code, "INVALID_REQUEST", "wrong suspension id must be invalid");
|
||||
const statusAfterMismatch = assertRpcSuccess(
|
||||
await request("gateway.suspend.status", { suspensionId }),
|
||||
"status after wrong resume",
|
||||
);
|
||||
assert.equal(statusAfterMismatch?.status, "ready", "wrong resume must preserve the lease");
|
||||
const resumed = assertRpcSuccess(
|
||||
await request("gateway.suspend.resume", { suspensionId }),
|
||||
"resume first lease",
|
||||
);
|
||||
assert.deepEqual(
|
||||
{ status: resumed?.status, resumed: resumed?.resumed },
|
||||
{ status: "running", resumed: true },
|
||||
"first resume must release the lease",
|
||||
);
|
||||
const repeatedResume = assertRpcSuccess(
|
||||
await request("gateway.suspend.resume", { suspensionId }),
|
||||
"repeat first resume",
|
||||
);
|
||||
assert.equal(repeatedResume?.resumed, false, "repeat resume must be idempotent");
|
||||
const recoveredHealth = await request("health");
|
||||
assert(hasGatewayHealthSummaryPayload(recoveredHealth), "health must return its full summary");
|
||||
} finally {
|
||||
ws.close();
|
||||
}
|
||||
}
|
||||
|
||||
export async function runGatewaySuspensionPreRestartClient(
|
||||
{
|
||||
statePath,
|
||||
@@ -354,43 +387,12 @@ export async function runGatewaySuspensionPreRestartClient(
|
||||
assert.equal(blockedAdminHealth.status, 503, "Admin health must return HTTP 503");
|
||||
assertGatewaySuspendingError(blockedAdminHealth.body);
|
||||
|
||||
const upgrade = await requestUpgradeRejection(url, remainingDeadlineMs(requestContext.deadline));
|
||||
assert.equal(upgrade.status, 503, "new websocket upgrade must return HTTP 503");
|
||||
assert.equal(
|
||||
upgrade.body,
|
||||
"Gateway websocket admission closed",
|
||||
"new websocket upgrade must return the canonical admission body",
|
||||
);
|
||||
|
||||
const wrongResume = await rpc("gateway.suspend.resume", {
|
||||
suspensionId: `${firstLease.suspensionId}-wrong`,
|
||||
await verifyPreparedSuspensionSocket({
|
||||
deadline: requestContext.deadline,
|
||||
suspensionId: firstLease.suspensionId,
|
||||
token,
|
||||
url,
|
||||
});
|
||||
assert.equal(wrongResume.status, 400, "wrong suspension id must return HTTP 400");
|
||||
assert.equal(
|
||||
wrongResume.body?.error?.code,
|
||||
"INVALID_REQUEST",
|
||||
"wrong suspension id must return INVALID_REQUEST",
|
||||
);
|
||||
const statusAfterMismatch = assertAdminSuccess(
|
||||
await rpc("gateway.suspend.status", { suspensionId: firstLease.suspensionId }),
|
||||
"status after wrong resume",
|
||||
);
|
||||
assert.equal(statusAfterMismatch?.status, "ready", "wrong resume must preserve the lease");
|
||||
|
||||
const resumed = assertAdminSuccess(
|
||||
await rpc("gateway.suspend.resume", { suspensionId: firstLease.suspensionId }),
|
||||
"resume first lease",
|
||||
);
|
||||
assert.deepEqual(
|
||||
{ status: resumed?.status, resumed: resumed?.resumed },
|
||||
{ status: "running", resumed: true },
|
||||
"first resume must release the lease",
|
||||
);
|
||||
const repeatedResume = assertAdminSuccess(
|
||||
await rpc("gateway.suspend.resume", { suspensionId: firstLease.suspensionId }),
|
||||
"repeat first resume",
|
||||
);
|
||||
assert.equal(repeatedResume?.resumed, false, "repeat resume must be idempotent");
|
||||
|
||||
assertHealthyProbes(
|
||||
await readProbe(requestContext, "/healthz"),
|
||||
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
runGatewayNetworkClient,
|
||||
runGatewaySuspensionPostRestartClient,
|
||||
runGatewaySuspensionPreRestartClient,
|
||||
verifyPreparedSuspensionSocket,
|
||||
} from "../../scripts/e2e/lib/gateway-network/client.mts";
|
||||
import { readGatewayNetworkClientConnectTimeoutMs } from "../../scripts/e2e/lib/gateway-network/limits.mts";
|
||||
import { onceFrame } from "../../scripts/e2e/lib/gateway-network/ws-frames.mts";
|
||||
@@ -498,4 +499,123 @@ describe("gateway network client", () => {
|
||||
}),
|
||||
).toThrow("identify gateway-draining");
|
||||
});
|
||||
function createPreparedSocketHarness(responses: GatewayFrame[]) {
|
||||
const frames = [...responses];
|
||||
const requests: Array<{ method: string; params: Record<string, unknown> }> = [];
|
||||
let closeCount = 0;
|
||||
const socket = {
|
||||
close: () => {
|
||||
closeCount += 1;
|
||||
},
|
||||
send: (payload: string) => {
|
||||
const frame = JSON.parse(payload) as {
|
||||
method: string;
|
||||
params: Record<string, unknown>;
|
||||
};
|
||||
requests.push({ method: frame.method, params: frame.params });
|
||||
},
|
||||
};
|
||||
return {
|
||||
get closeCount() {
|
||||
return closeCount;
|
||||
},
|
||||
requests,
|
||||
deps: {
|
||||
onceFrame: async (
|
||||
_ws: unknown,
|
||||
predicate: (frame: GatewayFrame) => boolean,
|
||||
_timeoutMs?: number,
|
||||
) => {
|
||||
const response = frames.shift();
|
||||
expect(response).toBeDefined();
|
||||
const frame = {
|
||||
type: "res",
|
||||
id: `s${requests.length}`,
|
||||
...response,
|
||||
};
|
||||
expect(predicate(frame)).toBe(true);
|
||||
return frame;
|
||||
},
|
||||
openSocket: async () => socket,
|
||||
protocolVersion: 1,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
it("uses one authenticated socket for the prepared suspension control lifecycle", async () => {
|
||||
const suspending = {
|
||||
ok: false,
|
||||
error: {
|
||||
code: "UNAVAILABLE",
|
||||
retryable: true,
|
||||
details: { reason: "gateway-suspending", phase: "prepared" },
|
||||
},
|
||||
};
|
||||
const harness = createPreparedSocketHarness([
|
||||
{ ok: true },
|
||||
{ ok: true, payload: { status: "ready" } },
|
||||
suspending,
|
||||
{ ok: false, error: { code: "INVALID_REQUEST" } },
|
||||
{ ok: true, payload: { status: "ready" } },
|
||||
{ ok: true, payload: { status: "running", resumed: true } },
|
||||
{ ok: true, payload: { status: "running", resumed: false } },
|
||||
healthResponse(),
|
||||
]);
|
||||
await verifyPreparedSuspensionSocket(
|
||||
{
|
||||
deadline: Date.now() + 1_000,
|
||||
suspensionId: "lease-1",
|
||||
token: "test-token",
|
||||
url: "ws://127.0.0.1:12345",
|
||||
},
|
||||
harness.deps,
|
||||
);
|
||||
expect(harness.requests).toEqual([
|
||||
{
|
||||
method: "connect",
|
||||
params: {
|
||||
minProtocol: 1,
|
||||
maxProtocol: 1,
|
||||
client: {
|
||||
id: "cli",
|
||||
displayName: "docker-net-e2e",
|
||||
version: "dev",
|
||||
platform: process.platform,
|
||||
mode: "cli",
|
||||
},
|
||||
caps: [],
|
||||
auth: { token: "test-token" },
|
||||
role: "operator",
|
||||
scopes: ["operator.admin"],
|
||||
},
|
||||
},
|
||||
{ method: "gateway.suspend.status", params: { suspensionId: "lease-1" } },
|
||||
{ method: "health", params: {} },
|
||||
{
|
||||
method: "gateway.suspend.resume",
|
||||
params: { suspensionId: "lease-1-wrong" },
|
||||
},
|
||||
{ method: "gateway.suspend.status", params: { suspensionId: "lease-1" } },
|
||||
{ method: "gateway.suspend.resume", params: { suspensionId: "lease-1" } },
|
||||
{ method: "gateway.suspend.resume", params: { suspensionId: "lease-1" } },
|
||||
{ method: "health", params: {} },
|
||||
]);
|
||||
expect(harness.closeCount).toBe(1);
|
||||
const invalidHarness = createPreparedSocketHarness([
|
||||
{ ok: true },
|
||||
{ ok: true, payload: { status: "running" } },
|
||||
]);
|
||||
await expect(
|
||||
verifyPreparedSuspensionSocket(
|
||||
{
|
||||
deadline: Date.now() + 1_000,
|
||||
suspensionId: "lease-1",
|
||||
token: "test-token",
|
||||
url: "ws://127.0.0.1:12345",
|
||||
},
|
||||
invalidHarness.deps,
|
||||
),
|
||||
).rejects.toThrow("prepared suspension must remain ready");
|
||||
expect(invalidHarness.closeCount).toBe(1);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user