diff --git a/extensions/slack/src/probe.test.ts b/extensions/slack/src/probe.test.ts index b9cb46d8f1c7..51401c15182e 100644 --- a/extensions/slack/src/probe.test.ts +++ b/extensions/slack/src/probe.test.ts @@ -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), + }); }); }); diff --git a/extensions/slack/src/probe.transport.test.ts b/extensions/slack/src/probe.transport.test.ts new file mode 100644 index 000000000000..efec55127ce8 --- /dev/null +++ b/extensions/slack/src/probe.transport.test.ts @@ -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; + close(): Promise; +}; + +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 { + const server = createServer(handler); + const sockets = new Set(); + server.on("connection", (socket) => { + sockets.add(socket); + socket.once("close", () => sockets.delete(socket)); + }); + await new Promise((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((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((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(); + } + }); +}); diff --git a/extensions/slack/src/probe.ts b/extensions/slack/src/probe.ts index 87310b1e7ecc..b6eb0e76a110 100644 --- a/extensions/slack/src/probe.ts +++ b/extensions/slack/src/probe.ts @@ -19,7 +19,13 @@ export async function probeSlack( timeoutMs = 2500, opts?: { accountId?: string | null; identity?: "bot" | "user" }, ): Promise { - 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 () => {