mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-12 21:53:00 -06:00
fix(gateway): log websocket handshake phase (#93402)
* fix(gateway): log websocket handshake phase * fix(gateway): clarify websocket handshake phases
This commit is contained in:
@@ -4,6 +4,7 @@ import type { ResolvedGatewayAuth } from "../auth.js";
|
||||
import { MAX_BUFFERED_BYTES } from "../server-constants.js";
|
||||
import {
|
||||
attachGatewayWsForTest,
|
||||
createGatewayWsTestLogger,
|
||||
createGatewayWsTestRequestContext,
|
||||
createGatewayWsTestSocket,
|
||||
createResolvedGatewayTokenAuth,
|
||||
@@ -47,18 +48,20 @@ async function connectTestWs(
|
||||
options?: Partial<Parameters<typeof attachGatewayWsConnectionHandler>[0]>;
|
||||
} = {},
|
||||
) {
|
||||
const logWsControl = createGatewayWsTestLogger();
|
||||
const connected = attachGatewayWsForTest({
|
||||
attach: attachGatewayWsConnectionHandler,
|
||||
clients: params.clients,
|
||||
headers: params.headers,
|
||||
host: params.host,
|
||||
options: params.options,
|
||||
options: { ...params.options, logWsControl: logWsControl as never },
|
||||
socket: params.socket,
|
||||
});
|
||||
await waitForLazyMessageHandler();
|
||||
|
||||
return {
|
||||
clients: connected.clients,
|
||||
logWsControl,
|
||||
socket: connected.socket,
|
||||
passed: firstAttachedHandlerParams(),
|
||||
};
|
||||
@@ -197,6 +200,76 @@ describe("attachGatewayWsConnectionHandler", () => {
|
||||
expect(socket.close).toHaveBeenCalledWith(1008, "slow consumer");
|
||||
});
|
||||
|
||||
it("keeps handshake phase advancement monotonic", async () => {
|
||||
const { socket, logWsControl, passed } = await connectTestWs();
|
||||
const handlerParams = passed as {
|
||||
advanceHandshakePhase: (phase: string) => void;
|
||||
};
|
||||
|
||||
handlerParams.advanceHandshakePhase("auth_credentials_received");
|
||||
handlerParams.advanceHandshakePhase("auth_validated");
|
||||
handlerParams.advanceHandshakePhase("auth_credentials_received");
|
||||
socket.emit("close", 1006, Buffer.from("client disappeared"));
|
||||
|
||||
const [message, context] = logWsControl.warn.mock.calls[0] as [string, { phase?: string }];
|
||||
expect(message).toContain("phase=auth_validated");
|
||||
expect(context).toMatchObject({ phase: "auth_validated" });
|
||||
});
|
||||
|
||||
it("includes the last completed handshake phase in pre-connect close logs", async () => {
|
||||
const { socket, logWsControl } = await connectTestWs();
|
||||
|
||||
socket.emit("close", 1006, Buffer.from("client disappeared"));
|
||||
|
||||
expect(logWsControl.warn).toHaveBeenCalled();
|
||||
const [message, context] = logWsControl.warn.mock.calls[0] as [string, { phase?: string }];
|
||||
expect(message).toContain("closed before connect");
|
||||
expect(message).toContain("phase=ws_upgrade_started");
|
||||
expect(context).toMatchObject({ phase: "ws_upgrade_started" });
|
||||
});
|
||||
|
||||
it("includes the last completed handshake phase on preauth timeout logs", async () => {
|
||||
vi.useFakeTimers();
|
||||
const { logWsControl } = await connectTestWs({
|
||||
options: { preauthHandshakeTimeoutMs: 100 },
|
||||
});
|
||||
|
||||
vi.advanceTimersByTime(150);
|
||||
|
||||
expect(logWsControl.warn).toHaveBeenCalledWith(expect.stringContaining("handshake timeout"));
|
||||
expect(logWsControl.warn).toHaveBeenCalledWith(
|
||||
expect.stringContaining("phase=ws_upgrade_started"),
|
||||
);
|
||||
});
|
||||
|
||||
it("omits handshake phase metadata after the connection is ready", async () => {
|
||||
const { socket, logWsControl, passed } = await connectTestWs();
|
||||
const handlerParams = passed as {
|
||||
advanceHandshakePhase: (phase: string) => void;
|
||||
setClient: (client: never) => boolean;
|
||||
setHandshakeState: (state: "pending" | "connected" | "failed") => void;
|
||||
};
|
||||
|
||||
handlerParams.advanceHandshakePhase("auth_credentials_received");
|
||||
handlerParams.advanceHandshakePhase("auth_validated");
|
||||
expect(
|
||||
handlerParams.setClient({
|
||||
socket,
|
||||
connect: { client: { id: "openclaw-control-ui", mode: "webchat" } },
|
||||
connId: "ready-client",
|
||||
usesSharedGatewayAuth: false,
|
||||
} as never),
|
||||
).toBe(true);
|
||||
handlerParams.setHandshakeState("connected");
|
||||
handlerParams.advanceHandshakePhase("session_attached");
|
||||
handlerParams.advanceHandshakePhase("hello_payload_prepared");
|
||||
handlerParams.advanceHandshakePhase("ready");
|
||||
|
||||
socket.emit("close", 1000, Buffer.from("done"));
|
||||
|
||||
expect(logWsControl.warn).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips node presence disconnects for stale reconnected sockets", async () => {
|
||||
const unregister = vi.fn(() => null);
|
||||
const { socket } = attachGatewayWsForTest({
|
||||
|
||||
@@ -44,7 +44,7 @@ import type {
|
||||
WsOriginCheckMetrics,
|
||||
} from "./ws-connection/message-handler.js";
|
||||
import { resolveSharedGatewaySessionGeneration } from "./ws-shared-generation.js";
|
||||
import type { GatewayWsClient } from "./ws-types.js";
|
||||
import { WS_HANDSHAKE_PHASES, type GatewayWsClient, type WsHandshakePhase } from "./ws-types.js";
|
||||
|
||||
type SubsystemLogger = ReturnType<typeof createSubsystemLogger>;
|
||||
|
||||
@@ -289,6 +289,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
|
||||
logWs("in", "open", { connId, remoteAddr, remotePort, localAddr, localPort, endpoint });
|
||||
let handshakeState: "pending" | "connected" | "failed" = "pending";
|
||||
let lastHandshakePhase: WsHandshakePhase = "tcp_accepted";
|
||||
let holdsPreauthBudget = true;
|
||||
let closeCause: string | undefined;
|
||||
let closeMeta: Record<string, unknown> = {};
|
||||
@@ -296,6 +297,12 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
let lastFrameMethod: string | undefined;
|
||||
let lastFrameId: string | undefined;
|
||||
|
||||
const advanceHandshakePhase = (next: WsHandshakePhase) => {
|
||||
if (WS_HANDSHAKE_PHASES.indexOf(next) > WS_HANDSHAKE_PHASES.indexOf(lastHandshakePhase)) {
|
||||
lastHandshakePhase = next;
|
||||
}
|
||||
};
|
||||
|
||||
const setCloseCause = (cause: string, meta?: Record<string, unknown>) => {
|
||||
if (!closeCause) {
|
||||
closeCause = cause;
|
||||
@@ -331,9 +338,10 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
setCloseCause("handshake-timeout", {
|
||||
handshakeMs: Date.now() - openedAt,
|
||||
endpoint,
|
||||
phase: lastHandshakePhase,
|
||||
});
|
||||
logWsControl.warn(
|
||||
`handshake timeout conn=${connId} peer=${endpoint ?? "n/a"} remote=${remoteAddr ?? "?"}`,
|
||||
`handshake timeout conn=${connId} peer=${endpoint ?? "n/a"} remote=${remoteAddr ?? "?"} phase=${lastHandshakePhase}`,
|
||||
);
|
||||
close();
|
||||
}
|
||||
@@ -390,6 +398,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
event: "connect.challenge",
|
||||
payload: { nonce: connectNonce, ts: Date.now() },
|
||||
});
|
||||
advanceHandshakePhase("ws_upgrade_started");
|
||||
|
||||
socket.once("error", (err) => {
|
||||
if (isWsPayloadLimitError(err)) {
|
||||
@@ -414,9 +423,11 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
const logHost = sanitizeLogValue(requestHost);
|
||||
const logUserAgent = sanitizeLogValue(requestUserAgent);
|
||||
const logReason = sanitizeLogValue(reason?.toString());
|
||||
const handshakeIncomplete = lastHandshakePhase !== "ready";
|
||||
const closeContext = {
|
||||
cause: closeCause,
|
||||
handshake: handshakeState,
|
||||
...(handshakeIncomplete ? { phase: lastHandshakePhase } : {}),
|
||||
durationMs,
|
||||
lastFrameType,
|
||||
lastFrameMethod,
|
||||
@@ -467,7 +478,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
? ` suppressed=${closeLogDecision.suppressedSinceLastLog}`
|
||||
: "";
|
||||
logFn(
|
||||
`closed before connect conn=${connId} peer=${endpoint ?? "n/a"} remote=${remoteAddr ?? "?"} fwd=${logForwardedFor || "n/a"} origin=${logOrigin || "n/a"} host=${logHost || "n/a"} ua=${logUserAgent || "n/a"} code=${code ?? "n/a"} reason=${logReason || "n/a"}${suppressedText}`,
|
||||
`closed before connect conn=${connId} peer=${endpoint ?? "n/a"} remote=${remoteAddr ?? "?"} fwd=${logForwardedFor || "n/a"} origin=${logOrigin || "n/a"} host=${logHost || "n/a"} ua=${logUserAgent || "n/a"} code=${code ?? "n/a"} reason=${logReason || "n/a"} phase=${lastHandshakePhase}${suppressedText}`,
|
||||
closeContext,
|
||||
);
|
||||
}
|
||||
@@ -506,6 +517,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
durationMs,
|
||||
cause: closeCause,
|
||||
handshake: handshakeState,
|
||||
...(handshakeIncomplete ? { phase: lastHandshakePhase } : {}),
|
||||
lastFrameType,
|
||||
lastFrameMethod,
|
||||
lastFrameId,
|
||||
@@ -567,6 +579,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
setHandshakeState: (next) => {
|
||||
handshakeState = next;
|
||||
},
|
||||
advanceHandshakePhase,
|
||||
setCloseCause,
|
||||
setLastFrameMeta,
|
||||
originCheckMetrics,
|
||||
|
||||
@@ -212,6 +212,7 @@ function attachGatewayHarness(options: {
|
||||
mode: "none",
|
||||
allowTailscale: false,
|
||||
};
|
||||
const advanceHandshakePhase = vi.fn();
|
||||
attachGatewayWsMessageHandler({
|
||||
socket,
|
||||
upgradeReq: {
|
||||
@@ -244,6 +245,7 @@ function attachGatewayHarness(options: {
|
||||
return true;
|
||||
},
|
||||
setHandshakeState: vi.fn(),
|
||||
advanceHandshakePhase,
|
||||
setCloseCause: options.setCloseCause ?? createSetCloseCauseMock(),
|
||||
setLastFrameMeta: vi.fn(),
|
||||
originCheckMetrics: { hostHeaderFallbackAccepted: 0 },
|
||||
@@ -256,6 +258,7 @@ function attachGatewayHarness(options: {
|
||||
}
|
||||
const sendMessage = onMessage;
|
||||
return {
|
||||
advanceHandshakePhase,
|
||||
socketSend,
|
||||
sendRequest: (id: string, method: string, params: Record<string, unknown> = {}) => {
|
||||
sendMessage(
|
||||
@@ -584,6 +587,46 @@ describe("attachGatewayWsMessageHandler post-connect health refresh", () => {
|
||||
expect(JSON.stringify(captured.events)).not.toContain("gateway-token");
|
||||
});
|
||||
|
||||
it("records credential and hello preparation phases during connect", async () => {
|
||||
const harness = attachGatewayHarness({
|
||||
connId: "conn-phases",
|
||||
connectNonce: "nonce-phases",
|
||||
resolvedAuth: {
|
||||
mode: "token",
|
||||
token: "gateway-token",
|
||||
allowTailscale: false,
|
||||
},
|
||||
});
|
||||
|
||||
harness.sendConnect("connect-phases", {
|
||||
minProtocol: PROTOCOL_VERSION,
|
||||
maxProtocol: PROTOCOL_VERSION,
|
||||
client: {
|
||||
id: "gateway-client",
|
||||
version: "dev",
|
||||
platform: "test",
|
||||
mode: "backend",
|
||||
},
|
||||
role: "operator",
|
||||
scopes: [],
|
||||
caps: [],
|
||||
auth: {
|
||||
token: "gateway-token",
|
||||
},
|
||||
});
|
||||
|
||||
await vi.waitFor(() => {
|
||||
expect(harness.socketSend).toHaveBeenCalled();
|
||||
});
|
||||
expect(harness.advanceHandshakePhase.mock.calls.map(([phase]) => phase)).toEqual([
|
||||
"auth_credentials_received",
|
||||
"auth_validated",
|
||||
"session_attached",
|
||||
"hello_payload_prepared",
|
||||
"ready",
|
||||
]);
|
||||
});
|
||||
|
||||
it("does not mark local backend self-pairing clients as approval runtimes", async () => {
|
||||
const refreshHealthSnapshot = vi.fn<GatewayRequestContext["refreshHealthSnapshot"]>(async () =>
|
||||
createHealthSummary(),
|
||||
|
||||
@@ -157,7 +157,7 @@ import {
|
||||
incrementPresenceVersion,
|
||||
} from "../health-state.js";
|
||||
import { resolveSharedGatewaySessionGeneration } from "../ws-shared-generation.js";
|
||||
import type { GatewayWsClient } from "../ws-types.js";
|
||||
import type { GatewayWsClient, WsHandshakePhase } from "../ws-types.js";
|
||||
import { resolveConnectAuthDecision, resolveConnectAuthState } from "./auth-context.js";
|
||||
import { formatGatewayAuthFailureMessage } from "./auth-messages.js";
|
||||
import {
|
||||
@@ -493,6 +493,7 @@ export type GatewayWsMessageHandlerParams = {
|
||||
getClient: () => GatewayWsClient | null;
|
||||
setClient: (next: GatewayWsClient) => boolean;
|
||||
setHandshakeState: (state: "pending" | "connected" | "failed") => void;
|
||||
advanceHandshakePhase: (phase: WsHandshakePhase) => void;
|
||||
setCloseCause: (cause: string, meta?: Record<string, unknown>) => void;
|
||||
setLastFrameMeta: (meta: { type?: string; method?: string; id?: string }) => void;
|
||||
originCheckMetrics: WsOriginCheckMetrics;
|
||||
@@ -538,6 +539,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
|
||||
getClient,
|
||||
setClient,
|
||||
setHandshakeState,
|
||||
advanceHandshakePhase,
|
||||
setCloseCause,
|
||||
setLastFrameMeta,
|
||||
originCheckMetrics,
|
||||
@@ -928,6 +930,14 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
|
||||
deviceRaw,
|
||||
});
|
||||
const device = controlUiAuthPolicy.device;
|
||||
const hasRawHandshakeCredentials =
|
||||
hasSharedAuth ||
|
||||
Boolean(connectParams.auth?.bootstrapToken) ||
|
||||
Boolean(connectParams.auth?.deviceToken) ||
|
||||
Boolean(device);
|
||||
if (hasRawHandshakeCredentials) {
|
||||
advanceHandshakePhase("auth_credentials_received");
|
||||
}
|
||||
const connectAuthState = await resolveConnectAuthState({
|
||||
resolvedAuth,
|
||||
connectAuth: connectParams.auth,
|
||||
@@ -1272,6 +1282,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
|
||||
rejectUnauthorized(authResult);
|
||||
return;
|
||||
}
|
||||
advanceHandshakePhase("auth_validated");
|
||||
const usesSharedGatewayAuth =
|
||||
authMethod === "token" || authMethod === "password" || authMethod === "trusted-proxy";
|
||||
const sharedGatewaySessionGeneration = usesSharedGatewayAuth
|
||||
@@ -2051,6 +2062,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
|
||||
return;
|
||||
}
|
||||
setHandshakeState("connected");
|
||||
advanceHandshakePhase("session_attached");
|
||||
logWs("in", "connect", {
|
||||
connId,
|
||||
client: connectParams.client.id,
|
||||
@@ -2190,6 +2202,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
|
||||
tickIntervalMs: TICK_INTERVAL_MS,
|
||||
},
|
||||
};
|
||||
advanceHandshakePhase("hello_payload_prepared");
|
||||
|
||||
let revokedBootstrapTokenRecord:
|
||||
| Awaited<ReturnType<typeof revokeDeviceBootstrapToken>>["record"]
|
||||
@@ -2257,6 +2270,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
|
||||
clientMode: connectParams.client.mode,
|
||||
deviceId: device?.id,
|
||||
});
|
||||
advanceHandshakePhase("ready");
|
||||
if (pendingNodePairingCleanup) {
|
||||
const context = buildRequestContext();
|
||||
const cleanupClaim = pendingNodePairingCleanup;
|
||||
|
||||
@@ -26,3 +26,15 @@ export type GatewayWsClient = PluginNodeCapabilityClient & {
|
||||
invalidated?: boolean;
|
||||
invalidatedReason?: string;
|
||||
};
|
||||
|
||||
export const WS_HANDSHAKE_PHASES = [
|
||||
"tcp_accepted",
|
||||
"ws_upgrade_started",
|
||||
"auth_credentials_received",
|
||||
"auth_validated",
|
||||
"session_attached",
|
||||
"hello_payload_prepared",
|
||||
"ready",
|
||||
] as const;
|
||||
|
||||
export type WsHandshakePhase = (typeof WS_HANDSHAKE_PHASES)[number];
|
||||
|
||||
Reference in New Issue
Block a user