From 96fada5483137ec044908ac0b23831776b98dd7a Mon Sep 17 00:00:00 2001 From: Frank Yang Date: Thu, 21 May 2026 15:40:05 +0800 Subject: [PATCH] fix(gateway): log slow node send diagnostics --- src/gateway/node-registry.test.ts | 27 ++++++++++++++++++++++++--- src/gateway/node-registry.ts | 7 +++++++ 2 files changed, 31 insertions(+), 3 deletions(-) diff --git a/src/gateway/node-registry.test.ts b/src/gateway/node-registry.test.ts index 1151a3611476..722dca8aa005 100644 --- a/src/gateway/node-registry.test.ts +++ b/src/gateway/node-registry.test.ts @@ -1,5 +1,6 @@ import { EventEmitter } from "node:events"; import { describe, expect, it, vi } from "vitest"; +import { onDiagnosticEvent, resetDiagnosticEventsForTest } from "../infra/diagnostic-events.js"; import { NodeRegistry, serializeEventPayload } from "./node-registry.js"; import { MAX_BUFFERED_BYTES } from "./server-constants.js"; import type { GatewayWsClient } from "./server/ws-types.js"; @@ -508,6 +509,9 @@ describe("gateway/node-registry", () => { }); it("rejects raw event sends when the node socket buffer is saturated", () => { + resetDiagnosticEventsForTest(); + const diagnosticEvents: unknown[] = []; + const stopDiagnostics = onDiagnosticEvent((event) => diagnosticEvents.push(event)); const registry = new NodeRegistry(); const socket = { bufferedAmount: MAX_BUFFERED_BYTES + 1, @@ -522,9 +526,26 @@ describe("gateway/node-registry", () => { ); const payload = serializeEventPayload({ foo: "bar" }); - expect(registry.sendEventRaw("node-1", "chat", payload)).toBe(false); - expect(socket.send).not.toHaveBeenCalled(); - expect(socket.close).toHaveBeenCalledWith(1008, "slow consumer"); + try { + expect(registry.sendEventRaw("node-1", "chat", payload)).toBe(false); + expect(socket.send).not.toHaveBeenCalled(); + expect(socket.close).toHaveBeenCalledWith(1008, "slow consumer"); + expect(diagnosticEvents).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + type: "payload.large", + action: "rejected", + surface: "gateway.ws.outbound_buffer", + bytes: MAX_BUFFERED_BYTES + 1, + limitBytes: MAX_BUFFERED_BYTES, + reason: "ws_send_buffer_close", + }), + ]), + ); + } finally { + stopDiagnostics(); + resetDiagnosticEventsForTest(); + } }); it("refreshes effective live surface within the declared surface", () => { diff --git a/src/gateway/node-registry.ts b/src/gateway/node-registry.ts index 23353d24b665..2d523d54f205 100644 --- a/src/gateway/node-registry.ts +++ b/src/gateway/node-registry.ts @@ -1,4 +1,5 @@ import { randomUUID } from "node:crypto"; +import { logRejectedLargePayload } from "../logging/diagnostic-payload.js"; import { MAX_BUFFERED_BYTES } from "./server-constants.js"; import type { GatewayWsClient } from "./server/ws-types.js"; @@ -710,6 +711,12 @@ export class NodeRegistry { if (!(node.client.socket.bufferedAmount > MAX_BUFFERED_BYTES)) { return false; } + logRejectedLargePayload({ + surface: "gateway.ws.outbound_buffer", + bytes: node.client.socket.bufferedAmount, + limitBytes: MAX_BUFFERED_BYTES, + reason: "ws_send_buffer_close", + }); try { node.client.socket.close(SLOW_CONSUMER_CLOSE_CODE, "slow consumer"); } catch {