From bf310355387dc77e82f547d0e9557f1cd3532035 Mon Sep 17 00:00:00 2001 From: Yzx <53250620+849261680@users.noreply.github.com> Date: Wed, 8 Jul 2026 11:31:51 +0800 Subject: [PATCH] fix(gateway): log websocket handshake phase (#93402) * fix(gateway): log websocket handshake phase * fix(gateway): clarify websocket handshake phases --- src/gateway/server/ws-connection.test.ts | 75 ++++++++++++++++++- src/gateway/server/ws-connection.ts | 19 ++++- ...essage-handler.post-connect-health.test.ts | 43 +++++++++++ .../server/ws-connection/message-handler.ts | 16 +++- src/gateway/server/ws-types.ts | 12 +++ 5 files changed, 160 insertions(+), 5 deletions(-) diff --git a/src/gateway/server/ws-connection.test.ts b/src/gateway/server/ws-connection.test.ts index f203579f75f6..859dfa5490eb 100644 --- a/src/gateway/server/ws-connection.test.ts +++ b/src/gateway/server/ws-connection.test.ts @@ -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[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({ diff --git a/src/gateway/server/ws-connection.ts b/src/gateway/server/ws-connection.ts index d515f5b9f65a..08602dcf9a71 100644 --- a/src/gateway/server/ws-connection.ts +++ b/src/gateway/server/ws-connection.ts @@ -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; @@ -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 = {}; @@ -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) => { 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, diff --git a/src/gateway/server/ws-connection/message-handler.post-connect-health.test.ts b/src/gateway/server/ws-connection/message-handler.post-connect-health.test.ts index d1c09a30da62..1db8c404f6b4 100644 --- a/src/gateway/server/ws-connection/message-handler.post-connect-health.test.ts +++ b/src/gateway/server/ws-connection/message-handler.post-connect-health.test.ts @@ -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 = {}) => { 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(async () => createHealthSummary(), diff --git a/src/gateway/server/ws-connection/message-handler.ts b/src/gateway/server/ws-connection/message-handler.ts index 6585469c76f6..1f3b322a4215 100644 --- a/src/gateway/server/ws-connection/message-handler.ts +++ b/src/gateway/server/ws-connection/message-handler.ts @@ -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) => 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>["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; diff --git a/src/gateway/server/ws-types.ts b/src/gateway/server/ws-types.ts index fc32bea0dccc..b8c314c4cf9f 100644 --- a/src/gateway/server/ws-types.ts +++ b/src/gateway/server/ws-types.ts @@ -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];