diff --git a/extensions/slack/src/monitor/presence-monitor.test.ts b/extensions/slack/src/monitor/presence-monitor.test.ts index 06d25d08300d..84fd8fef7c92 100644 --- a/extensions/slack/src/monitor/presence-monitor.test.ts +++ b/extensions/slack/src/monitor/presence-monitor.test.ts @@ -1,7 +1,12 @@ +import { WebAPIRateLimitedError } from "@slack/web-api"; import type { PluginStateSyncKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime"; import { describe, expect, it, vi } from "vitest"; import type { PreparedSlackMessage } from "./message-handler/types.js"; -import { createSlackPresenceMonitor, hasSlackPresenceEventsEnabled } from "./presence-monitor.js"; +import { + createSlackPresenceMonitor, + hasSlackPresenceEventsEnabled, + SLACK_PRESENCE_REQUEST_TIMEOUT_MS, +} from "./presence-monitor.js"; const AUTO_MAX_PARTICIPANTS = 8; @@ -298,6 +303,120 @@ describe("Slack presence monitor", () => { expect(enqueue).not.toHaveBeenCalled(); }); + it("times out a stalled presence request and polls the next user", async () => { + vi.useFakeTimers(); + let resolveStalled!: (value: { presence: string }) => void; + const stalled = new Promise<{ presence: string }>((resolve) => { + resolveStalled = resolve; + }); + let polling: Promise | undefined; + try { + const getPresence = vi + .fn() + .mockReturnValueOnce(stalled) + .mockResolvedValueOnce({ presence: "away" }); + const monitor = createSlackPresenceMonitor({ + accountId: "default", + accountConfig: { mode: "auto" }, + client: { getPresence } as never, + cooldownStore: createCooldownStore(), + enqueue: vi.fn(() => true), + wake: vi.fn(), + }); + monitor.observe(createPrepared({ userId: "U1", channelId: "D1" })); + monitor.observe(createPrepared({ userId: "U2", channelId: "D2" })); + + polling = monitor.pollOnce(); + let pollSettled = false; + void polling.then(() => { + pollSettled = true; + }); + await vi.advanceTimersByTimeAsync(SLACK_PRESENCE_REQUEST_TIMEOUT_MS); + expect(pollSettled).toBe(true); + await polling; + + expect(getPresence).toHaveBeenNthCalledWith(1, { user: "U1" }); + expect(getPresence).toHaveBeenNthCalledWith(2, { user: "U2" }); + } finally { + resolveStalled({ presence: "away" }); + await polling; + vi.useRealTimers(); + } + }); + + it("honors Slack Retry-After without skipping the unpolled page", async () => { + let now = 1_000; + const getPresence = vi + .fn() + .mockRejectedValueOnce(new WebAPIRateLimitedError(120)) + .mockResolvedValue({ presence: "away" }); + const monitor = createSlackPresenceMonitor({ + accountId: "default", + accountConfig: { mode: "on" }, + client: { getPresence } as never, + cooldownStore: createCooldownStore(), + enqueue: vi.fn(() => true), + wake: vi.fn(), + nowMs: () => now, + }); + for (let index = 1; index <= 46; index += 1) { + monitor.observe(createPrepared({ userId: `U${String(index).padStart(4, "0")}` })); + } + + await monitor.pollOnce(); + expect(getPresence).toHaveBeenCalledExactlyOnceWith({ user: "U0001" }); + + now += 119_999; + await monitor.pollOnce(); + expect(getPresence).toHaveBeenCalledTimes(1); + + now += 1; + await monitor.pollOnce(); + expect(getPresence).toHaveBeenNthCalledWith(2, { user: "U0001" }); + expect(getPresence).toHaveBeenNthCalledWith(3, { user: "U0002" }); + expect(getPresence).toHaveBeenCalledTimes(46); + + await monitor.pollOnce(); + expect(getPresence).toHaveBeenNthCalledWith(47, { user: "U0046" }); + }); + + it("bounds stop while a presence request is stalled", async () => { + vi.useFakeTimers(); + let resolveStalled!: (value: { presence: string }) => void; + const stalled = new Promise<{ presence: string }>((resolve) => { + resolveStalled = resolve; + }); + let polling: Promise | undefined; + try { + const getPresence = vi.fn(() => stalled); + const monitor = createSlackPresenceMonitor({ + accountId: "default", + accountConfig: { mode: "auto" }, + client: { getPresence } as never, + cooldownStore: createCooldownStore(), + enqueue: vi.fn(() => true), + wake: vi.fn(), + }); + monitor.observe(createPrepared({ userId: "U1" })); + + polling = monitor.pollOnce(); + const stopping = monitor.stop(); + let stopSettled = false; + void stopping.then(() => { + stopSettled = true; + }); + await vi.advanceTimersByTimeAsync(SLACK_PRESENCE_REQUEST_TIMEOUT_MS); + expect(stopSettled).toBe(true); + await Promise.all([polling, stopping]); + + expect(getPresence).toHaveBeenCalledOnce(); + } finally { + resolveStalled({ presence: "away" }); + await polling; + vi.useRealTimers(); + } + }); + it("quiesces an in-flight poll before stop returns", async () => { let resolveActive!: (value: { presence: string }) => void; const active = new Promise<{ presence: string }>((resolve) => { diff --git a/extensions/slack/src/monitor/presence-monitor.ts b/extensions/slack/src/monitor/presence-monitor.ts index 4a843db1b178..ec0f5d462fce 100644 --- a/extensions/slack/src/monitor/presence-monitor.ts +++ b/extensions/slack/src/monitor/presence-monitor.ts @@ -1,12 +1,14 @@ // Slack plugin module polls selected participants and routes away-to-active transitions. -import type { WebClient } from "@slack/web-api"; +import { type WebClient, WebAPIRateLimitedError } from "@slack/web-api"; import type { SlackAccountConfig } from "openclaw/plugin-sdk/config-contracts"; import { requestHeartbeat } from "openclaw/plugin-sdk/heartbeat-runtime"; import type { PluginStateSyncKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime"; import { enqueueSystemEvent } from "openclaw/plugin-sdk/system-event-runtime"; +import { withTimeout } from "openclaw/plugin-sdk/text-utility-runtime"; import type { PreparedSlackMessage } from "./message-handler/types.js"; export const SLACK_PRESENCE_GREETING_COOLDOWN_MS = 8 * 60 * 60 * 1000; +export const SLACK_PRESENCE_REQUEST_TIMEOUT_MS = 30_000; const SLACK_PRESENCE_POLL_INTERVAL_MS = 60_000; const SLACK_PRESENCE_AUTO_MAX_PARTICIPANTS = 8; const SLACK_PRESENCE_TARGET_TTL_MS = 24 * 60 * 60 * 1000; @@ -142,6 +144,7 @@ export function createSlackPresenceMonitor(params: { let pollOffset = 0; let timer: NodeJS.Timeout | undefined; let activePoll: Promise | undefined; + let rateLimitedUntilMs = 0; let stopped = false; const pruneTargets = (now: number) => { @@ -249,6 +252,10 @@ export function createSlackPresenceMonitor(params: { const performPoll = async () => { const now = nowMs(); + if (rateLimitedUntilMs > now) { + return; + } + rateLimitedUntilMs = 0; pruneTargets(now); const candidates = Array.from( new Set( @@ -265,16 +272,23 @@ export function createSlackPresenceMonitor(params: { { length: count }, (_, index) => candidates[(pollOffset + index) % candidates.length], ).filter((userId): userId is string => Boolean(userId)); - pollOffset = (pollOffset + count) % candidates.length; for (const userId of selected) { if (stopped) { return; } + let consumed = false; try { - const response = await params.client.getPresence({ user: userId }); + const response = await withTimeout( + params.client.getPresence({ user: userId }), + SLACK_PRESENCE_REQUEST_TIMEOUT_MS, + { + message: `Slack presence request timed out after ${SLACK_PRESENCE_REQUEST_TIMEOUT_MS}ms`, + }, + ); if (stopped) { return; } + consumed = true; const next = response.presence === "active" || response.presence === "away" ? response.presence @@ -291,7 +305,20 @@ export function createSlackPresenceMonitor(params: { if (stopped) { return; } + if (err instanceof WebAPIRateLimitedError) { + rateLimitedUntilMs = Math.max( + rateLimitedUntilMs, + nowMs() + Math.max(0, err.retryAfter) * 1_000, + ); + params.error?.(`slack presence polling rate limited; retrying after ${err.retryAfter}s`); + return; + } + consumed = true; params.error?.(`slack presence poll failed for user ${userId}: ${String(err)}`); + } finally { + if (consumed) { + pollOffset = (pollOffset + 1) % candidates.length; + } } } }; diff --git a/extensions/slack/src/monitor/provider.auth-test-token.test.ts b/extensions/slack/src/monitor/provider.auth-test-token.test.ts index b207cd2f8b04..4fda13233635 100644 --- a/extensions/slack/src/monitor/provider.auth-test-token.test.ts +++ b/extensions/slack/src/monitor/provider.auth-test-token.test.ts @@ -1,6 +1,8 @@ // Slack tests cover auth.test token handling during provider boot. import { createServer } from "node:http"; import type { AddressInfo } from "node:net"; +import type { OpenKeyedStoreOptions } from "openclaw/plugin-sdk/plugin-state-runtime"; +import { createPluginStateSyncKeyedStoreForTests } from "openclaw/plugin-sdk/plugin-state-test-runtime"; import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { disposeSlackTestRuntime, @@ -13,6 +15,7 @@ import { stopSlackMonitor, useRealSlackStartupAuthClientOnce, } from "../monitor.test-helpers.js"; +import { getSlackRuntime } from "../runtime.js"; const { monitorSlackProvider } = await import("./provider.js"); @@ -341,6 +344,87 @@ describe("auth.test boot call", () => { }); }); +describe("presence polling transport", () => { + it("aborts a stalled presence request when the provider stops", async () => { + const events: string[] = []; + for (const key of PROXY_ENV_KEYS) { + vi.stubEnv(key, ""); + } + const server = await startStalledSlackApiServer(events); + vi.stubEnv("SLACK_API_URL", server.apiUrl); + resetSlackTestState({ + channels: { + slack: { + dm: { enabled: true }, + dmPolicy: "open", + allowFrom: ["*"], + groupPolicy: "open", + presenceEvents: { mode: "on" }, + }, + }, + }); + getSlackRuntime().state.openSyncKeyedStore = (options: OpenKeyedStoreOptions) => + createPluginStateSyncKeyedStoreForTests("slack", { + ...options, + env: options.env ?? process.env, + }); + getSlackTestState().replyMock.mockResolvedValue({ text: "ok" }); + + const nativeSetInterval = globalThis.setInterval; + let triggerPresencePoll: (() => void) | undefined; + const intervalSpy = vi.spyOn(globalThis, "setInterval").mockImplementation((( + handler: (...args: unknown[]) => void, + timeout?: number, + ...args: unknown[] + ) => { + if (timeout === 60_000 && !triggerPresencePoll) { + triggerPresencePoll = () => handler(...args); + return nativeSetInterval(() => undefined, 60 * 60 * 1_000); + } + return nativeSetInterval(handler, timeout, ...args); + }) as typeof setInterval); + + const monitor = startSlackMonitor(monitorSlackProvider); + try { + const handler = await getSlackHandlerOrThrow("message"); + await handler({ + event: { + type: "message", + user: "U_STALLED", + text: "hello", + ts: "100.000", + channel: "D_STALLED", + channel_type: "im", + }, + context: { botUserId: "bot-user" }, + body: {}, + }); + expect(triggerPresencePoll).toBeTypeOf("function"); + triggerPresencePoll?.(); + await vi.waitFor(() => expect(server.requestCount).toBe(1), { timeout: 1_000 }); + + const startedAt = Date.now(); + monitor.controller.abort(); + const outcome = await Promise.race([ + monitor.run.then(() => "settled" as const), + new Promise<"timed-out">((resolve) => { + setTimeout(() => resolve("timed-out"), 2_000); + }), + ]); + + expect(outcome).toBe("settled"); + expect(Date.now() - startedAt).toBeLessThan(2_000); + await vi.waitFor(() => expect(events).toContain("socket-closed"), { timeout: 1_000 }); + expect(server.requestUrl).toBe("/api/users.getPresence"); + } finally { + intervalSpy.mockRestore(); + monitor.controller.abort(); + await server.close(); + await monitor.run; + } + }); +}); + describe("user identity provider transport", () => { const userSocketConfig = () => ({ channels: { diff --git a/extensions/slack/src/monitor/provider.ts b/extensions/slack/src/monitor/provider.ts index 69687ead8194..8b519b22cfbb 100644 --- a/extensions/slack/src/monitor/provider.ts +++ b/extensions/slack/src/monitor/provider.ts @@ -1,5 +1,6 @@ // Slack provider module implements model/runtime integration. import type { IncomingMessage, ServerResponse } from "node:http"; +import { type FetchFunction, WebClient } from "@slack/web-api"; import { addAllowlistUserEntriesFromConfigEntry, buildAllowlistResolutionSummary, @@ -33,7 +34,11 @@ import { resolveSlackAccountDmPolicy, } from "../accounts.js"; import { isSlackAnyNativeApprovalClientEnabled } from "../approval-native-gates.js"; -import { resolveSlackProxyDispatcher, resolveSlackWebClientOptions } from "../client-options.js"; +import { + resolveSlackLookupClientOptions, + resolveSlackProxyDispatcher, + resolveSlackWebClientOptions, +} from "../client-options.js"; import { createSlackStartupAuthClient } from "../client.js"; import { normalizeSlackWebhookPath, registerSlackHttpHandler } from "../http/index.js"; import { SLACK_TEXT_LIMIT } from "../limits.js"; @@ -66,7 +71,11 @@ import { registerSlackMonitorEvents } from "./events.js"; import { createSlackDurableIngress } from "./ingress.js"; import { createSlackMessageHandler } from "./message-handler.js"; import { openSlackPresenceCooldownStore } from "./presence-cooldown-store.js"; -import { createSlackPresenceMonitor, hasSlackPresenceEventsEnabled } from "./presence-monitor.js"; +import { + createSlackPresenceMonitor, + hasSlackPresenceEventsEnabled, + SLACK_PRESENCE_REQUEST_TIMEOUT_MS, +} from "./presence-monitor.js"; import { createSlackBoltApp, formatSlackChannelResolved, @@ -92,6 +101,17 @@ import type { MonitorSlackOpts } from "./types.js"; let slackBoltInterop: SlackBoltResolvedExports | undefined; +function withSlackPresenceLifecycleSignal( + fetchImpl: FetchFunction, + lifecycleSignal: AbortSignal, +): FetchFunction { + return async (input, init) => + await fetchImpl(input, { + ...init, + signal: init?.signal ? AbortSignal.any([init.signal, lifecycleSignal]) : lifecycleSignal, + }); +} + async function getSlackBoltInterop(): Promise { if (!slackBoltInterop) { const slackBoltModule = await import("@slack/bolt"); @@ -614,17 +634,34 @@ export async function monitorSlackProvider(opts: MonitorSlackOpts = {}) { account: slackCfg.presenceEvents, channels: slackCfg.channels, }); - const presenceMonitor = + const presenceRequestAbort = installationIdentity.kind !== "enterprise" && presenceEventsEnabled - ? createSlackPresenceMonitor({ - accountId: account.accountId, - accountConfig: slackCfg.presenceEvents, - client: app.client.users, - cooldownStore: openSlackPresenceCooldownStore(), - log: runtime.log, - error: runtime.error, - }) + ? new AbortController() : undefined; + const presenceClient = + presenceRequestAbort === undefined + ? undefined + : (() => { + const options = resolveSlackLookupClientOptions( + { ...clientOptions, timeout: SLACK_PRESENCE_REQUEST_TIMEOUT_MS }, + slackDispatcher, + ); + options.fetch = withSlackPresenceLifecycleSignal( + options.fetch ?? globalThis.fetch, + presenceRequestAbort.signal, + ); + return new WebClient(token, options).users; + })(); + const presenceMonitor = presenceClient + ? createSlackPresenceMonitor({ + accountId: account.accountId, + accountConfig: slackCfg.presenceEvents, + client: presenceClient, + cooldownStore: openSlackPresenceCooldownStore(), + log: runtime.log, + error: runtime.error, + }) + : undefined; if (installationIdentity.kind === "enterprise" && presenceEventsEnabled) { runtime.log?.(warn("slack presence events are unavailable for Enterprise Grid org installs")); } @@ -945,6 +982,7 @@ export async function monitorSlackProvider(opts: MonitorSlackOpts = {}) { } } } finally { + presenceRequestAbort?.abort(); await presenceMonitor?.stop(); if (slackMode === "relay") { setSlackDefaultSendIdentity(account.accountId, undefined);