From 95bbd117ef2e1d349c2458354cac2ef90ee97bb3 Mon Sep 17 00:00:00 2001 From: Calin Laurentiu Ilie Date: Thu, 13 Aug 2026 03:57:04 +0200 Subject: [PATCH] fix(slack): prevent duplicate Socket Mode connections after reconnect errors (#122624) * fix(codex): read canonical transcript session targets (#1) * test(slack): reproduce reconnect timer surviving shutdown * fix(slack): keep reconnects within one socket lifecycle * test(slack): exercise native reconnect over loopback * test(slack): satisfy reconnect integration checks --- .../slack/src/monitor/provider-support.ts | 6 +- .../src/monitor/provider.interop.test.ts | 134 +++++++++++++++++- .../src/monitor/provider.reconnect.test.ts | 5 +- .../slack/src/monitor/reconnect-policy.ts | 5 +- 4 files changed, 142 insertions(+), 8 deletions(-) diff --git a/extensions/slack/src/monitor/provider-support.ts b/extensions/slack/src/monitor/provider-support.ts index 7424ef2ba950..19342284080b 100644 --- a/extensions/slack/src/monitor/provider-support.ts +++ b/extensions/slack/src/monitor/provider-support.ts @@ -95,7 +95,8 @@ function installSlackNativeReconnectFailureObserver(receiver: unknown) { `Before trying to reconnect, this client will wait for ${delayMs} milliseconds`, ); return new Promise((resolve, reject) => { - setTimeout(() => { + const reconnectTimer = setTimeout(() => { + Reflect.set(this, "reconnectionTimer", undefined); if (Reflect.get(this, "shuttingDown")) { logger?.debug?.("Client shutting down, will not attempt reconnect."); resolve(undefined); @@ -112,6 +113,9 @@ function installSlackNativeReconnectFailureObserver(receiver: unknown) { reject(toErrorObject(error, "Non-Error rejection")); }); }, delayMs); + // SocketModeClient.disconnect() clears this field. Keep the patched + // scheduler on the SDK's lifecycle so a stopped app cannot reconnect. + Reflect.set(this, "reconnectionTimer", reconnectTimer); }); }, ); diff --git a/extensions/slack/src/monitor/provider.interop.test.ts b/extensions/slack/src/monitor/provider.interop.test.ts index a75d7926fbb6..d6a8aee389b3 100644 --- a/extensions/slack/src/monitor/provider.interop.test.ts +++ b/extensions/slack/src/monitor/provider.interop.test.ts @@ -1,6 +1,12 @@ // Slack tests cover provider.interop plugin behavior. import { describe, expect, it, vi } from "vitest"; -import { createSlackBoltApp, resolveSlackBoltInterop } from "./provider-support.js"; +import { WebSocketServer } from "ws"; +import { + createSlackBoltApp, + gracefulStopSlackApp, + resolveSlackBoltInterop, + startSlackSocketAndWaitForDisconnect, +} from "./provider-support.js"; describe("resolveSlackBoltInterop", () => { function FakeApp() {} @@ -356,6 +362,132 @@ describe("createSlackBoltApp", () => { ]); }); + it("cancels a pending native reconnect when the app is stopped and started again", async () => { + vi.useFakeTimers(); + try { + const slackBoltModule = await import("@slack/bolt"); + const { app, receiver } = createSlackBoltApp({ + interop: resolveSlackBoltInterop({ + defaultImport: slackBoltModule.default, + namespaceImport: slackBoltModule, + }), + slackMode: "socket", + token: "xoxb-test", + appToken: "xapp-test", + slackWebhookPath: "/slack/events", + clientOptions: {}, + }); + if (!receiver || typeof receiver !== "object") { + throw new Error("expected a Socket Mode receiver"); + } + const client = Reflect.get(receiver, "client"); + if (!client || typeof client !== "object") { + throw new Error("expected a Socket Mode client"); + } + const start = vi.fn(async () => { + Reflect.set(client, "shuttingDown", false); + }); + Reflect.set(client, "start", start); + const delayReconnectAttempt = Reflect.get(client, "delayReconnectAttempt"); + if (typeof delayReconnectAttempt !== "function") { + throw new Error("expected a native reconnect scheduler"); + } + + void delayReconnectAttempt.call(client, start); + await gracefulStopSlackApp(app); + await app.start(); + await vi.advanceTimersByTimeAsync(15_000); + + expect(start).toHaveBeenCalledTimes(1); + await gracefulStopSlackApp(app); + } finally { + vi.useRealTimers(); + } + }); + + it("recovers a transient error and close through one real SDK socket lifecycle", async () => { + const socketServer = new WebSocketServer({ port: 0 }); + await new Promise((resolve) => { + socketServer.once("listening", resolve); + }); + const address = socketServer.address(); + if (!address || typeof address === "string") { + throw new Error("expected a TCP Socket Mode test server"); + } + let connectionAttempts = 0; + let peakActiveConnections = 0; + socketServer.on("connection", (socket) => { + connectionAttempts += 1; + peakActiveConnections = Math.max(peakActiveConnections, socketServer.clients.size); + socket.send(JSON.stringify({ type: "hello", num_connections: socketServer.clients.size })); + }); + + const slackBoltModule = await import("@slack/bolt"); + const { app, receiver } = createSlackBoltApp({ + interop: resolveSlackBoltInterop({ + defaultImport: slackBoltModule.default, + namespaceImport: slackBoltModule, + }), + slackMode: "socket", + token: "xoxb-test", + appToken: "xapp-test", + slackWebhookPath: "/slack/events", + clientOptions: { + fetch: async () => + new Response( + JSON.stringify({ + ok: true, + url: `ws://127.0.0.1:${address.port}`, + }), + { headers: { "content-type": "application/json" } }, + ), + }, + }); + if (!receiver || typeof receiver !== "object") { + throw new Error("expected a Socket Mode receiver"); + } + const client = Reflect.get(receiver, "client"); + if (!client || typeof client !== "object") { + throw new Error("expected a Socket Mode client"); + } + Reflect.set(client, "clientPingTimeoutMS", 20); + const appStart = vi.spyOn(app, "start"); + const abortController = new AbortController(); + const lifecycle = startSlackSocketAndWaitForDisconnect({ + app, + abortSignal: abortController.signal, + }); + let lifecycleSettled = false; + const lifecycleOutcome = lifecycle.then((value) => { + lifecycleSettled = true; + return value; + }); + + try { + await vi.waitFor(() => expect(socketServer.clients.size).toBe(1)); + Reflect.get(client, "emit").call(client, "error", new Error("transient transport error")); + for (const socket of socketServer.clients) { + socket.terminate(); + } + await vi.waitFor(() => expect(connectionAttempts).toBe(2)); + await vi.waitFor(() => expect(socketServer.clients.size).toBe(1)); + + expect(appStart).toHaveBeenCalledTimes(1); + expect(peakActiveConnections).toBe(1); + expect(lifecycleSettled).toBe(false); + } finally { + abortController.abort(); + await lifecycleOutcome; + await gracefulStopSlackApp(app); + for (const socket of socketServer.clients) { + socket.terminate(); + } + await new Promise((resolve, reject) => { + socketServer.close((error) => (error ? reject(error) : resolve())); + }); + } + }); + it("uses Slack's fixed Socket Mode receiver policy", () => { const clientOptions = { teamId: "T1" }; const { receiver } = createSlackBoltApp({ diff --git a/extensions/slack/src/monitor/provider.reconnect.test.ts b/extensions/slack/src/monitor/provider.reconnect.test.ts index 5849a0b7fcf0..e246187f186d 100644 --- a/extensions/slack/src/monitor/provider.reconnect.test.ts +++ b/extensions/slack/src/monitor/provider.reconnect.test.ts @@ -231,15 +231,16 @@ describe("slack socket reconnect helpers", () => { await expect(waiter).resolves.toEqual({ event: "disconnect" }); }); - it("resolves disconnect waiter on socket error event", async () => { + it("leaves transient socket errors to the native reconnect lifecycle", async () => { const client = new FakeEmitter(); const app = { receiver: { client } }; const err = new Error("dns down"); const waiter = waitForSlackSocketDisconnect(app as never); client.emit("error", err); + client.emit("disconnected"); - await expect(waiter).resolves.toEqual({ event: "error", error: err }); + await expect(waiter).resolves.toEqual({ event: "disconnect" }); }); it("installs the disconnect waiter before socket start completes", async () => { diff --git a/extensions/slack/src/monitor/reconnect-policy.ts b/extensions/slack/src/monitor/reconnect-policy.ts index 059241922f11..24f24c428d4e 100644 --- a/extensions/slack/src/monitor/reconnect-policy.ts +++ b/extensions/slack/src/monitor/reconnect-policy.ts @@ -13,7 +13,7 @@ export const SLACK_SOCKET_RECONNECT_POLICY = { jitter: 0.25, } as const; -type SlackSocketDisconnectEvent = "disconnect" | "unable_to_socket_mode_start" | "error"; +type SlackSocketDisconnectEvent = "disconnect" | "unable_to_socket_mode_start"; type EmitterLike = { on: (event: string, listener: (...args: unknown[]) => void) => unknown; @@ -132,13 +132,11 @@ export function waitForSlackSocketDisconnect( const disconnectListener = () => resolveOnce({ event: "disconnect" }); const startFailListener = (error?: unknown) => resolveOnce({ event: "unable_to_socket_mode_start", error }); - const errorListener = (error: unknown) => resolveOnce({ event: "error", error }); const abortListener = () => resolveOnce({ event: "disconnect" }); const cleanup = () => { emitter.off("disconnected", disconnectListener); emitter.off("unable_to_socket_mode_start", startFailListener); - emitter.off("error", errorListener); abortSignal?.removeEventListener("abort", abortListener); }; @@ -149,7 +147,6 @@ export function waitForSlackSocketDisconnect( emitter.on("disconnected", disconnectListener); emitter.on("unable_to_socket_mode_start", startFailListener); - emitter.on("error", errorListener); abortSignal?.addEventListener("abort", abortListener, { once: true }); }); }