mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-27 12:56:01 -06:00
fix(slack): abort health probes within their deadline (#107835)
Preserve the current dedicated Slack read-client architecture and stop probe retries or HTTP 429 delays from surviving the shared health deadline. Verify stalled sockets, dropped-request background retries, immediate rate-limit rejection, and trickling response bodies with real loopback transport. Fixes #106565. Co-authored-by: Peter Steinberger <steipete@gmail.com> Co-authored-by: zhangqueping <3436352+zhangqueping@users.noreply.github.com>
This commit is contained in:
@@ -39,7 +39,11 @@ describe("probeSlack", () => {
|
||||
bot: { id: "U123", name: "openclaw-bot" },
|
||||
team: { id: "T123", name: "OpenClaw" },
|
||||
});
|
||||
expect(createSlackReadClientMock).toHaveBeenCalledWith("xoxb-test", { timeout: 2500 });
|
||||
expect(createSlackReadClientMock).toHaveBeenCalledWith("xoxb-test", {
|
||||
rejectRateLimitedCalls: true,
|
||||
retryConfig: { retries: 0 },
|
||||
timeout: 2500,
|
||||
});
|
||||
});
|
||||
|
||||
it("warns when auth.test looks like a user token in the bot token slot", async () => {
|
||||
@@ -111,5 +115,36 @@ describe("probeSlack", () => {
|
||||
expect(result.elapsedMs).toBe(35);
|
||||
expect(result.bot).toStrictEqual({ id: undefined, name: undefined });
|
||||
expect(result.team).toStrictEqual({ id: undefined, name: undefined });
|
||||
expect(createSlackReadClientMock).toHaveBeenCalledWith("xoxb-test", {
|
||||
rejectRateLimitedCalls: true,
|
||||
retryConfig: { retries: 0 },
|
||||
timeout: 2500,
|
||||
});
|
||||
});
|
||||
|
||||
it("passes a custom probe deadline to Slack's abortable read transport", async () => {
|
||||
authTestMock.mockResolvedValue({ ok: true });
|
||||
|
||||
await expect(probeSlack("xoxb-test", 175)).resolves.toMatchObject({ ok: true });
|
||||
|
||||
expect(createSlackReadClientMock).toHaveBeenCalledWith("xoxb-test", {
|
||||
rejectRateLimitedCalls: true,
|
||||
retryConfig: { retries: 0 },
|
||||
timeout: 175,
|
||||
});
|
||||
});
|
||||
|
||||
it("keeps the normal health result when the Slack read transport aborts", async () => {
|
||||
authTestMock.mockRejectedValue(
|
||||
Object.assign(new Error("The operation was aborted due to timeout"), {
|
||||
name: "TimeoutError",
|
||||
}),
|
||||
);
|
||||
|
||||
await expect(probeSlack("xoxb-test", 175)).resolves.toMatchObject({
|
||||
ok: false,
|
||||
status: null,
|
||||
error: expect.any(String),
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,191 @@
|
||||
// Prove Slack probe deadlines against the real SDK and loopback HTTP transport.
|
||||
import { createServer, type RequestListener } from "node:http";
|
||||
import type { Socket } from "node:net";
|
||||
import { afterEach, describe, expect, it } from "vitest";
|
||||
import { probeSlack } from "./probe.js";
|
||||
|
||||
const TEST_ENV_KEYS = [
|
||||
"SLACK_API_URL",
|
||||
"HTTPS_PROXY",
|
||||
"HTTP_PROXY",
|
||||
"ALL_PROXY",
|
||||
"https_proxy",
|
||||
"http_proxy",
|
||||
"all_proxy",
|
||||
"NO_PROXY",
|
||||
"no_proxy",
|
||||
"OPENCLAW_PROXY_ACTIVE",
|
||||
"OPENCLAW_PROXY_CA_FILE",
|
||||
] as const;
|
||||
|
||||
const originalEnv = Object.fromEntries(
|
||||
TEST_ENV_KEYS.map((key) => [key, process.env[key]]),
|
||||
) as Record<(typeof TEST_ENV_KEYS)[number], string | undefined>;
|
||||
|
||||
type TestServer = {
|
||||
apiUrl: string;
|
||||
sockets: Set<Socket>;
|
||||
close(): Promise<void>;
|
||||
};
|
||||
|
||||
function clearSlackTransportEnv(): void {
|
||||
for (const key of TEST_ENV_KEYS) {
|
||||
delete process.env[key];
|
||||
}
|
||||
}
|
||||
|
||||
function restoreSlackTransportEnv(): void {
|
||||
for (const key of TEST_ENV_KEYS) {
|
||||
const original = originalEnv[key];
|
||||
if (original === undefined) {
|
||||
delete process.env[key];
|
||||
} else {
|
||||
process.env[key] = original;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function startSlackTransportServer(handler: RequestListener): Promise<TestServer> {
|
||||
const server = createServer(handler);
|
||||
const sockets = new Set<Socket>();
|
||||
server.on("connection", (socket) => {
|
||||
sockets.add(socket);
|
||||
socket.once("close", () => sockets.delete(socket));
|
||||
});
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.once("error", reject);
|
||||
server.listen(0, "127.0.0.1", () => {
|
||||
server.off("error", reject);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
const address = server.address();
|
||||
if (!address || typeof address === "string") {
|
||||
server.close();
|
||||
throw new Error("Slack probe test server did not bind a TCP address");
|
||||
}
|
||||
return {
|
||||
apiUrl: `http://127.0.0.1:${address.port}/api/`,
|
||||
sockets,
|
||||
close: async () => {
|
||||
for (const socket of sockets) {
|
||||
socket.destroy();
|
||||
}
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.close((error) => {
|
||||
if (error) {
|
||||
reject(error);
|
||||
} else {
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
});
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
restoreSlackTransportEnv();
|
||||
});
|
||||
|
||||
describe("probeSlack real network deadlines", () => {
|
||||
it("aborts a stalled Slack request and closes its only socket", async () => {
|
||||
clearSlackTransportEnv();
|
||||
let requests = 0;
|
||||
const server = await startSlackTransportServer((request) => {
|
||||
requests += 1;
|
||||
request.resume();
|
||||
});
|
||||
try {
|
||||
process.env.SLACK_API_URL = server.apiUrl;
|
||||
|
||||
await expect(probeSlack("probe-fixture", 100)).resolves.toMatchObject({ ok: false });
|
||||
|
||||
expect(requests).toBe(1);
|
||||
await expect.poll(() => server.sockets.size, { timeout: 400 }).toBe(0);
|
||||
} finally {
|
||||
await server.close();
|
||||
}
|
||||
});
|
||||
|
||||
it("does not retry a dropped request after the probe has already returned", async () => {
|
||||
clearSlackTransportEnv();
|
||||
let requests = 0;
|
||||
const server = await startSlackTransportServer((request, response) => {
|
||||
requests += 1;
|
||||
request.resume();
|
||||
response.destroy();
|
||||
});
|
||||
try {
|
||||
process.env.SLACK_API_URL = server.apiUrl;
|
||||
|
||||
await expect(probeSlack("probe-fixture", 100)).resolves.toMatchObject({ ok: false });
|
||||
|
||||
// The default Slack read retry is randomized between 500 and 1,000 ms.
|
||||
// Keep the endpoint alive long enough to catch work after the probe ends.
|
||||
await new Promise<void>((resolve) => {
|
||||
setTimeout(resolve, 1_200);
|
||||
});
|
||||
expect(requests).toBe(1);
|
||||
await expect.poll(() => server.sockets.size, { timeout: 400 }).toBe(0);
|
||||
} finally {
|
||||
await server.close();
|
||||
}
|
||||
});
|
||||
|
||||
it("rejects a rate limit without waiting or retrying", async () => {
|
||||
clearSlackTransportEnv();
|
||||
let requests = 0;
|
||||
const server = await startSlackTransportServer((request, response) => {
|
||||
requests += 1;
|
||||
request.resume();
|
||||
response.writeHead(429, {
|
||||
"content-type": "application/json",
|
||||
"retry-after": "2",
|
||||
});
|
||||
response.end(`${JSON.stringify({ ok: false, error: "ratelimited" })}\n`);
|
||||
});
|
||||
try {
|
||||
process.env.SLACK_API_URL = server.apiUrl;
|
||||
const startedAt = performance.now();
|
||||
|
||||
await expect(probeSlack("probe-fixture", 1_000)).resolves.toMatchObject({ ok: false });
|
||||
|
||||
expect(performance.now() - startedAt).toBeLessThan(500);
|
||||
expect(requests).toBe(1);
|
||||
} finally {
|
||||
await server.close();
|
||||
}
|
||||
});
|
||||
|
||||
it("aborts response-body trickling at the absolute probe deadline", async () => {
|
||||
clearSlackTransportEnv();
|
||||
let requests = 0;
|
||||
const server = await startSlackTransportServer((request, response) => {
|
||||
requests += 1;
|
||||
request.resume();
|
||||
response.writeHead(200, { "content-type": "application/json" });
|
||||
const trickle = setInterval(() => response.write(" "), 25);
|
||||
const finish = setTimeout(() => {
|
||||
clearInterval(trickle);
|
||||
response.end(`${JSON.stringify({ ok: true })}\n`);
|
||||
}, 750);
|
||||
response.once("close", () => {
|
||||
clearInterval(trickle);
|
||||
clearTimeout(finish);
|
||||
});
|
||||
});
|
||||
try {
|
||||
process.env.SLACK_API_URL = server.apiUrl;
|
||||
const startedAt = performance.now();
|
||||
|
||||
await expect(probeSlack("probe-fixture", 100)).resolves.toMatchObject({ ok: false });
|
||||
|
||||
expect(performance.now() - startedAt).toBeLessThan(500);
|
||||
expect(requests).toBe(1);
|
||||
await expect.poll(() => server.sockets.size, { timeout: 400 }).toBe(0);
|
||||
} finally {
|
||||
await server.close();
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -19,7 +19,13 @@ export async function probeSlack(
|
||||
timeoutMs = 2500,
|
||||
opts?: { accountId?: string | null; identity?: "bot" | "user" },
|
||||
): Promise<SlackProbe> {
|
||||
const client = createSlackReadClient(token, { timeout: timeoutMs });
|
||||
// The probe owns a single absolute deadline: abort its fetch and never let
|
||||
// retries or Slack's 429 queue outlive the shared health-check result.
|
||||
const client = createSlackReadClient(token, {
|
||||
rejectRateLimitedCalls: true,
|
||||
retryConfig: { retries: 0 },
|
||||
timeout: timeoutMs,
|
||||
});
|
||||
return await runChannelProbe(
|
||||
timeoutMs,
|
||||
async () => {
|
||||
|
||||
Reference in New Issue
Block a user