mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-12 21:53:00 -06:00
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
This commit is contained in:
committed by
GitHub
parent
d00e4ee324
commit
95bbd117ef
@@ -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);
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
@@ -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<void>((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<void>((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({
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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 });
|
||||
});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user