Skip to content

Commit bf31035

Browse files
authored
fix(gateway): log websocket handshake phase (#93402)
* fix(gateway): log websocket handshake phase * fix(gateway): clarify websocket handshake phases
1 parent 082bd45 commit bf31035

5 files changed

Lines changed: 160 additions & 5 deletions

File tree

src/gateway/server/ws-connection.test.ts

Lines changed: 74 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import type { ResolvedGatewayAuth } from "../auth.js";
44
import { MAX_BUFFERED_BYTES } from "../server-constants.js";
55
import {
66
attachGatewayWsForTest,
7+
createGatewayWsTestLogger,
78
createGatewayWsTestRequestContext,
89
createGatewayWsTestSocket,
910
createResolvedGatewayTokenAuth,
@@ -47,18 +48,20 @@ async function connectTestWs(
4748
options?: Partial<Parameters<typeof attachGatewayWsConnectionHandler>[0]>;
4849
} = {},
4950
) {
51+
const logWsControl = createGatewayWsTestLogger();
5052
const connected = attachGatewayWsForTest({
5153
attach: attachGatewayWsConnectionHandler,
5254
clients: params.clients,
5355
headers: params.headers,
5456
host: params.host,
55-
options: params.options,
57+
options: { ...params.options, logWsControl: logWsControl as never },
5658
socket: params.socket,
5759
});
5860
await waitForLazyMessageHandler();
5961

6062
return {
6163
clients: connected.clients,
64+
logWsControl,
6265
socket: connected.socket,
6366
passed: firstAttachedHandlerParams(),
6467
};
@@ -197,6 +200,76 @@ describe("attachGatewayWsConnectionHandler", () => {
197200
expect(socket.close).toHaveBeenCalledWith(1008, "slow consumer");
198201
});
199202

203+
it("keeps handshake phase advancement monotonic", async () => {
204+
const { socket, logWsControl, passed } = await connectTestWs();
205+
const handlerParams = passed as {
206+
advanceHandshakePhase: (phase: string) => void;
207+
};
208+
209+
handlerParams.advanceHandshakePhase("auth_credentials_received");
210+
handlerParams.advanceHandshakePhase("auth_validated");
211+
handlerParams.advanceHandshakePhase("auth_credentials_received");
212+
socket.emit("close", 1006, Buffer.from("client disappeared"));
213+
214+
const [message, context] = logWsControl.warn.mock.calls[0] as [string, { phase?: string }];
215+
expect(message).toContain("phase=auth_validated");
216+
expect(context).toMatchObject({ phase: "auth_validated" });
217+
});
218+
219+
it("includes the last completed handshake phase in pre-connect close logs", async () => {
220+
const { socket, logWsControl } = await connectTestWs();
221+
222+
socket.emit("close", 1006, Buffer.from("client disappeared"));
223+
224+
expect(logWsControl.warn).toHaveBeenCalled();
225+
const [message, context] = logWsControl.warn.mock.calls[0] as [string, { phase?: string }];
226+
expect(message).toContain("closed before connect");
227+
expect(message).toContain("phase=ws_upgrade_started");
228+
expect(context).toMatchObject({ phase: "ws_upgrade_started" });
229+
});
230+
231+
it("includes the last completed handshake phase on preauth timeout logs", async () => {
232+
vi.useFakeTimers();
233+
const { logWsControl } = await connectTestWs({
234+
options: { preauthHandshakeTimeoutMs: 100 },
235+
});
236+
237+
vi.advanceTimersByTime(150);
238+
239+
expect(logWsControl.warn).toHaveBeenCalledWith(expect.stringContaining("handshake timeout"));
240+
expect(logWsControl.warn).toHaveBeenCalledWith(
241+
expect.stringContaining("phase=ws_upgrade_started"),
242+
);
243+
});
244+
245+
it("omits handshake phase metadata after the connection is ready", async () => {
246+
const { socket, logWsControl, passed } = await connectTestWs();
247+
const handlerParams = passed as {
248+
advanceHandshakePhase: (phase: string) => void;
249+
setClient: (client: never) => boolean;
250+
setHandshakeState: (state: "pending" | "connected" | "failed") => void;
251+
};
252+
253+
handlerParams.advanceHandshakePhase("auth_credentials_received");
254+
handlerParams.advanceHandshakePhase("auth_validated");
255+
expect(
256+
handlerParams.setClient({
257+
socket,
258+
connect: { client: { id: "openclaw-control-ui", mode: "webchat" } },
259+
connId: "ready-client",
260+
usesSharedGatewayAuth: false,
261+
} as never),
262+
).toBe(true);
263+
handlerParams.setHandshakeState("connected");
264+
handlerParams.advanceHandshakePhase("session_attached");
265+
handlerParams.advanceHandshakePhase("hello_payload_prepared");
266+
handlerParams.advanceHandshakePhase("ready");
267+
268+
socket.emit("close", 1000, Buffer.from("done"));
269+
270+
expect(logWsControl.warn).not.toHaveBeenCalled();
271+
});
272+
200273
it("skips node presence disconnects for stale reconnected sockets", async () => {
201274
const unregister = vi.fn(() => null);
202275
const { socket } = attachGatewayWsForTest({

src/gateway/server/ws-connection.ts

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ import type {
4444
WsOriginCheckMetrics,
4545
} from "./ws-connection/message-handler.js";
4646
import { resolveSharedGatewaySessionGeneration } from "./ws-shared-generation.js";
47-
import type { GatewayWsClient } from "./ws-types.js";
47+
import { WS_HANDSHAKE_PHASES, type GatewayWsClient, type WsHandshakePhase } from "./ws-types.js";
4848

4949
type SubsystemLogger = ReturnType<typeof createSubsystemLogger>;
5050

@@ -289,13 +289,20 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
289289

290290
logWs("in", "open", { connId, remoteAddr, remotePort, localAddr, localPort, endpoint });
291291
let handshakeState: "pending" | "connected" | "failed" = "pending";
292+
let lastHandshakePhase: WsHandshakePhase = "tcp_accepted";
292293
let holdsPreauthBudget = true;
293294
let closeCause: string | undefined;
294295
let closeMeta: Record<string, unknown> = {};
295296
let lastFrameType: string | undefined;
296297
let lastFrameMethod: string | undefined;
297298
let lastFrameId: string | undefined;
298299

300+
const advanceHandshakePhase = (next: WsHandshakePhase) => {
301+
if (WS_HANDSHAKE_PHASES.indexOf(next) > WS_HANDSHAKE_PHASES.indexOf(lastHandshakePhase)) {
302+
lastHandshakePhase = next;
303+
}
304+
};
305+
299306
const setCloseCause = (cause: string, meta?: Record<string, unknown>) => {
300307
if (!closeCause) {
301308
closeCause = cause;
@@ -331,9 +338,10 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
331338
setCloseCause("handshake-timeout", {
332339
handshakeMs: Date.now() - openedAt,
333340
endpoint,
341+
phase: lastHandshakePhase,
334342
});
335343
logWsControl.warn(
336-
`handshake timeout conn=${connId} peer=${endpoint ?? "n/a"} remote=${remoteAddr ?? "?"}`,
344+
`handshake timeout conn=${connId} peer=${endpoint ?? "n/a"} remote=${remoteAddr ?? "?"} phase=${lastHandshakePhase}`,
337345
);
338346
close();
339347
}
@@ -390,6 +398,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
390398
event: "connect.challenge",
391399
payload: { nonce: connectNonce, ts: Date.now() },
392400
});
401+
advanceHandshakePhase("ws_upgrade_started");
393402

394403
socket.once("error", (err) => {
395404
if (isWsPayloadLimitError(err)) {
@@ -414,9 +423,11 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
414423
const logHost = sanitizeLogValue(requestHost);
415424
const logUserAgent = sanitizeLogValue(requestUserAgent);
416425
const logReason = sanitizeLogValue(reason?.toString());
426+
const handshakeIncomplete = lastHandshakePhase !== "ready";
417427
const closeContext = {
418428
cause: closeCause,
419429
handshake: handshakeState,
430+
...(handshakeIncomplete ? { phase: lastHandshakePhase } : {}),
420431
durationMs,
421432
lastFrameType,
422433
lastFrameMethod,
@@ -467,7 +478,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
467478
? ` suppressed=${closeLogDecision.suppressedSinceLastLog}`
468479
: "";
469480
logFn(
470-
`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}`,
481+
`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}`,
471482
closeContext,
472483
);
473484
}
@@ -506,6 +517,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
506517
durationMs,
507518
cause: closeCause,
508519
handshake: handshakeState,
520+
...(handshakeIncomplete ? { phase: lastHandshakePhase } : {}),
509521
lastFrameType,
510522
lastFrameMethod,
511523
lastFrameId,
@@ -567,6 +579,7 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
567579
setHandshakeState: (next) => {
568580
handshakeState = next;
569581
},
582+
advanceHandshakePhase,
570583
setCloseCause,
571584
setLastFrameMeta,
572585
originCheckMetrics,

src/gateway/server/ws-connection/message-handler.post-connect-health.test.ts

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -212,6 +212,7 @@ function attachGatewayHarness(options: {
212212
mode: "none",
213213
allowTailscale: false,
214214
};
215+
const advanceHandshakePhase = vi.fn();
215216
attachGatewayWsMessageHandler({
216217
socket,
217218
upgradeReq: {
@@ -244,6 +245,7 @@ function attachGatewayHarness(options: {
244245
return true;
245246
},
246247
setHandshakeState: vi.fn(),
248+
advanceHandshakePhase,
247249
setCloseCause: options.setCloseCause ?? createSetCloseCauseMock(),
248250
setLastFrameMeta: vi.fn(),
249251
originCheckMetrics: { hostHeaderFallbackAccepted: 0 },
@@ -256,6 +258,7 @@ function attachGatewayHarness(options: {
256258
}
257259
const sendMessage = onMessage;
258260
return {
261+
advanceHandshakePhase,
259262
socketSend,
260263
sendRequest: (id: string, method: string, params: Record<string, unknown> = {}) => {
261264
sendMessage(
@@ -584,6 +587,46 @@ describe("attachGatewayWsMessageHandler post-connect health refresh", () => {
584587
expect(JSON.stringify(captured.events)).not.toContain("gateway-token");
585588
});
586589

590+
it("records credential and hello preparation phases during connect", async () => {
591+
const harness = attachGatewayHarness({
592+
connId: "conn-phases",
593+
connectNonce: "nonce-phases",
594+
resolvedAuth: {
595+
mode: "token",
596+
token: "gateway-token",
597+
allowTailscale: false,
598+
},
599+
});
600+
601+
harness.sendConnect("connect-phases", {
602+
minProtocol: PROTOCOL_VERSION,
603+
maxProtocol: PROTOCOL_VERSION,
604+
client: {
605+
id: "gateway-client",
606+
version: "dev",
607+
platform: "test",
608+
mode: "backend",
609+
},
610+
role: "operator",
611+
scopes: [],
612+
caps: [],
613+
auth: {
614+
token: "gateway-token",
615+
},
616+
});
617+
618+
await vi.waitFor(() => {
619+
expect(harness.socketSend).toHaveBeenCalled();
620+
});
621+
expect(harness.advanceHandshakePhase.mock.calls.map(([phase]) => phase)).toEqual([
622+
"auth_credentials_received",
623+
"auth_validated",
624+
"session_attached",
625+
"hello_payload_prepared",
626+
"ready",
627+
]);
628+
});
629+
587630
it("does not mark local backend self-pairing clients as approval runtimes", async () => {
588631
const refreshHealthSnapshot = vi.fn<GatewayRequestContext["refreshHealthSnapshot"]>(async () =>
589632
createHealthSummary(),

src/gateway/server/ws-connection/message-handler.ts

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,7 @@ import {
157157
incrementPresenceVersion,
158158
} from "../health-state.js";
159159
import { resolveSharedGatewaySessionGeneration } from "../ws-shared-generation.js";
160-
import type { GatewayWsClient } from "../ws-types.js";
160+
import type { GatewayWsClient, WsHandshakePhase } from "../ws-types.js";
161161
import { resolveConnectAuthDecision, resolveConnectAuthState } from "./auth-context.js";
162162
import { formatGatewayAuthFailureMessage } from "./auth-messages.js";
163163
import {
@@ -493,6 +493,7 @@ export type GatewayWsMessageHandlerParams = {
493493
getClient: () => GatewayWsClient | null;
494494
setClient: (next: GatewayWsClient) => boolean;
495495
setHandshakeState: (state: "pending" | "connected" | "failed") => void;
496+
advanceHandshakePhase: (phase: WsHandshakePhase) => void;
496497
setCloseCause: (cause: string, meta?: Record<string, unknown>) => void;
497498
setLastFrameMeta: (meta: { type?: string; method?: string; id?: string }) => void;
498499
originCheckMetrics: WsOriginCheckMetrics;
@@ -538,6 +539,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
538539
getClient,
539540
setClient,
540541
setHandshakeState,
542+
advanceHandshakePhase,
541543
setCloseCause,
542544
setLastFrameMeta,
543545
originCheckMetrics,
@@ -928,6 +930,14 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
928930
deviceRaw,
929931
});
930932
const device = controlUiAuthPolicy.device;
933+
const hasRawHandshakeCredentials =
934+
hasSharedAuth ||
935+
Boolean(connectParams.auth?.bootstrapToken) ||
936+
Boolean(connectParams.auth?.deviceToken) ||
937+
Boolean(device);
938+
if (hasRawHandshakeCredentials) {
939+
advanceHandshakePhase("auth_credentials_received");
940+
}
931941
const connectAuthState = await resolveConnectAuthState({
932942
resolvedAuth,
933943
connectAuth: connectParams.auth,
@@ -1272,6 +1282,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
12721282
rejectUnauthorized(authResult);
12731283
return;
12741284
}
1285+
advanceHandshakePhase("auth_validated");
12751286
const usesSharedGatewayAuth =
12761287
authMethod === "token" || authMethod === "password" || authMethod === "trusted-proxy";
12771288
const sharedGatewaySessionGeneration = usesSharedGatewayAuth
@@ -2051,6 +2062,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
20512062
return;
20522063
}
20532064
setHandshakeState("connected");
2065+
advanceHandshakePhase("session_attached");
20542066
logWs("in", "connect", {
20552067
connId,
20562068
client: connectParams.client.id,
@@ -2190,6 +2202,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
21902202
tickIntervalMs: TICK_INTERVAL_MS,
21912203
},
21922204
};
2205+
advanceHandshakePhase("hello_payload_prepared");
21932206

21942207
let revokedBootstrapTokenRecord:
21952208
| Awaited<ReturnType<typeof revokeDeviceBootstrapToken>>["record"]
@@ -2257,6 +2270,7 @@ export function attachGatewayWsMessageHandler(params: GatewayWsMessageHandlerPar
22572270
clientMode: connectParams.client.mode,
22582271
deviceId: device?.id,
22592272
});
2273+
advanceHandshakePhase("ready");
22602274
if (pendingNodePairingCleanup) {
22612275
const context = buildRequestContext();
22622276
const cleanupClaim = pendingNodePairingCleanup;

src/gateway/server/ws-types.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,3 +26,15 @@ export type GatewayWsClient = PluginNodeCapabilityClient & {
2626
invalidated?: boolean;
2727
invalidatedReason?: string;
2828
};
29+
30+
export const WS_HANDSHAKE_PHASES = [
31+
"tcp_accepted",
32+
"ws_upgrade_started",
33+
"auth_credentials_received",
34+
"auth_validated",
35+
"session_attached",
36+
"hello_payload_prepared",
37+
"ready",
38+
] as const;
39+
40+
export type WsHandshakePhase = (typeof WS_HANDSHAKE_PHASES)[number];

0 commit comments

Comments
 (0)