diff --git a/apps/server/src/routes/slack/commands.test.ts b/apps/server/src/routes/slack/commands.test.ts index 81e97e97..f417389d 100644 --- a/apps/server/src/routes/slack/commands.test.ts +++ b/apps/server/src/routes/slack/commands.test.ts @@ -54,7 +54,7 @@ describe("handleSlackCommand (members only)", () => { slackTestState.resolveWorkspace = (teamId: string) => teamId === "T_KNOWN" ? Promise.resolve({ - workspace: { id: 1 }, + workspace: { id: 1, limits: { "slack-agent": true } }, botToken: "xoxb-test", botUserId: "UBOT", }) diff --git a/apps/server/src/routes/slack/commands.ts b/apps/server/src/routes/slack/commands.ts index 91966fdd..d5762db5 100644 --- a/apps/server/src/routes/slack/commands.ts +++ b/apps/server/src/routes/slack/commands.ts @@ -16,7 +16,12 @@ import { LINK_ACCOUNT_TEXT, } from "./blocks"; import type { SlackConfig, SlackEnv } from "./config"; -import { linkAccountUrl, requireSlackMember } from "./require-slack-member"; +import { + linkAccountUrl, + planRequiredMessage, + requireSlackMember, + slackAgentAllowed, +} from "./require-slack-member"; import { resolvePageFromUrl } from "./resolve-page"; import { resolveWorkspace } from "./workspace-resolver"; @@ -144,6 +149,9 @@ async function runMemberCommand( text: "openstatus isn't connected to this Slack workspace. Connect it from the openstatus dashboard.", }; } + if (!slackAgentAllowed(resolved.workspace)) { + return planRequiredMessage(config); + } const actor = await requireSlackMember({ workspace: resolved.workspace, teamId: command.team_id, diff --git a/apps/server/src/routes/slack/handler.test.ts b/apps/server/src/routes/slack/handler.test.ts index 457418bc..d325579f 100644 --- a/apps/server/src/routes/slack/handler.test.ts +++ b/apps/server/src/routes/slack/handler.test.ts @@ -1,5 +1,8 @@ import crypto from "node:crypto"; +import { db, eq } from "@openstatus/db"; +import { integration } from "@openstatus/db/src/schema"; +import { createTestWorkspace } from "@openstatus/db/src/test/factories"; import { beforeEach, describe, expect, test } from "@openstatus/test-utils"; import { Hono } from "hono"; @@ -19,7 +22,7 @@ import { looksLikeUncardedDraft, toolTaskTitle, } from "./handler"; -import { abortTurn, endTurn, startTurn } from "./running-turns"; +import { abortTurn, broadcastStop, endTurn, startTurn } from "./running-turns"; import { verifySlackSignature } from "./verify"; function createTestApp() { @@ -104,7 +107,7 @@ function resetSlackTestState() { name: "Test Workspace", slug: "test", plan: "free", - limits: {}, + limits: { "slack-agent": true }, }, botToken: "xoxb-test", botUserId: "UBOT", @@ -1047,21 +1050,27 @@ describe("streaming the agent's answer", () => { ]); }); - test("falls back to the placeholder when the workspace has no streaming", async () => { + test("falls back to the thread status when the workspace has no streaming", async () => { slackTestState.chatStreamEnabled = false; streamTurn(async () => {}); await mention("3"); await new Promise((r) => setTimeout(r, 100)); - // No agent session and no stream leaves the oldest path: a "Thinking..." - // message posted up front and overwritten with the answer. - const placeholder = slackTestState.calls.find( + // No agent session and no stream leaves the thread status as the only + // loading indicator; the answer is posted as a message of its own. + expect( + slackTestState.calls.some( + (m) => m.method === "assistant.threads.setStatus", + ), + ).toBe(true); + const posts = slackTestState.calls.filter( (m) => m.method === "postMessage", ); - expect(placeholder?.args.text).toContain("Thinking..."); - const answer = slackTestState.calls.find((m) => m.method === "update"); - expect(answer?.args.text).toBe("All five monitors are healthy."); + expect(posts.map((m) => m.args.text)).toEqual([ + "All five monitors are healthy.", + ]); + expect(slackTestState.calls.some((m) => m.method === "update")).toBe(false); expect(slackTestState.calls.some((m) => m.method === "chatStream")).toBe( false, ); @@ -1783,3 +1792,143 @@ describe("members only", () => { expect(view.blocks.some((b) => b.type === "actions")).toBe(true); }); }); + +describe("hardening", () => { + const app = createTestApp(); + const redisStore = (globalThis as Record) + .__testRedisStore as Map; + + beforeEach(resetSlackTestState); + + test("a workspace without the Slack agent gets a plan notice", async () => { + slackTestState.resolveWorkspace = () => + Promise.resolve({ + workspace: { + id: 1, + name: "Free", + slug: "free", + plan: "free", + limits: {}, + }, + botToken: "xoxb-test", + botUserId: "UBOT", + }); + const user = `U_FREE_${crypto.randomUUID()}`; + await signAndPost(app, { + type: "event_callback", + team_id: "T_FREE", + event_id: `evt_free_${Date.now()}`, + event: { + type: "app_mention", + text: "<@UBOT> hi", + user, + channel: "C1", + channel_type: "channel", + ts: `${Date.now()}.61`, + }, + }); + const notice = await waitForCall("postEphemeral"); + expect(String(notice?.args.text)).toContain("/settings/billing"); + // The plan notice must not use up the account-link card's window. + expect(redisStore.has(`slack:linkcard:T_FREE:${user}`)).toBe(false); + }); + + test("the same event id is processed once, across instances", async () => { + const eventId = `evt_dup_${Date.now()}`; + const body = { + type: "event_callback", + team_id: "T_KNOWN", + event_id: eventId, + event: { + type: "app_mention", + text: "<@UBOT> status?", + user: "U1", + channel: "C1", + channel_type: "channel", + ts: `${Date.now()}.62`, + }, + }; + await signAndPost(app, body); + await signAndPost(app, body); + expect(redisStore.has(`slack:event:${eventId}`)).toBe(true); + await new Promise((r) => setTimeout(r, 150)); + const thinking = slackTestState.calls.filter( + (c) => + c.method === "chatStream" || + (c.method === "postMessage" && + String(c.args.text).includes("Thinking")), + ); + expect(thinking).toHaveLength(1); + }); + + test("a stop from another instance aborts the running turn", async () => { + const turn = startTurn("C_REMOTE", "7.7"); + await new Promise((r) => setTimeout(r, 10)); + await broadcastStop("C_REMOTE", "7.7"); + await new Promise((r) => setTimeout(r, 1_200)); + expect(turn.signal.aborted).toBe(true); + endTurn("C_REMOTE", "7.7", turn); + }); + + test("an earlier stop does not abort a newer turn", async () => { + await broadcastStop("C_OLD", "8.8"); + const turn = startTurn("C_OLD", "8.8"); + await new Promise((r) => setTimeout(r, 1_200)); + expect(turn.signal.aborted).toBe(false); + endTurn("C_OLD", "8.8", turn); + }); + + async function seedSlackInstall() { + const teamId = `T_REVOKED_${crypto.randomUUID()}`; + const { workspace } = await createTestWorkspace(); + const row = await db + .insert(integration) + .values({ + name: "slack-agent", + workspaceId: workspace.id, + externalId: teamId, + data: {}, + }) + .returning() + .get(); + return { teamId, id: row.id }; + } + + async function revokeTokens( + teamId: string, + tokens: Record, + ) { + const res = await signAndPost(app, { + type: "event_callback", + team_id: teamId, + event_id: `evt_revoked_${crypto.randomUUID()}`, + event: { type: "tokens_revoked", tokens }, + }); + expect(res.status).toBe(200); + await settleBackgroundTasks(); + } + + function findIntegration(id: number) { + return db.select().from(integration).where(eq(integration.id, id)).get(); + } + + test("tokens_revoked without a bot token leaves the install alone", async () => { + const install = await seedSlackInstall(); + try { + await revokeTokens(install.teamId, { oauth: ["U1"] }); + expect(await findIntegration(install.id)).toBeDefined(); + } finally { + await db.delete(integration).where(eq(integration.id, install.id)); + } + }); + + test("tokens_revoked with a bot token uninstalls", async () => { + const install = await seedSlackInstall(); + try { + await revokeTokens(install.teamId, { bot: ["UBOT"] }); + expect(await findIntegration(install.id)).toBeUndefined(); + } finally { + await db.delete(integration).where(eq(integration.id, install.id)); + } + }); +}); diff --git a/apps/server/src/routes/slack/handler.ts b/apps/server/src/routes/slack/handler.ts index c5cfb531..ac8fc4d2 100644 --- a/apps/server/src/routes/slack/handler.ts +++ b/apps/server/src/routes/slack/handler.ts @@ -1,10 +1,11 @@ import { getLogger } from "@logtape/logtape"; -import { and, db, eq, isNull, sql } from "@openstatus/db"; -import { integration, pageSubscriber } from "@openstatus/db/src/schema"; +import { uninstallSlackTeam } from "@openstatus/services/integration"; import { WebClient } from "@slack/web-api"; import type { Context } from "hono"; import { z } from "zod"; +import { redis } from "@/libs/clients"; + import { type AgentEvents, runAgent } from "./agent"; import { greetOnce, setAssistantStatus, setSessionStatus } from "./assistant"; import { runInBackground } from "./background"; @@ -41,11 +42,14 @@ import { } from "./registry-runner"; import { claimLinkCardWindow, + claimPlanNoticeWindow, linkAccountUrl, + planRequiredMessage, releaseLinkCardWindow, requireSlackMember, + slackAgentAllowed, } from "./require-slack-member"; -import { abortTurn, endTurn, startTurn } from "./running-turns"; +import { abortTurn, broadcastStop, endTurn, startTurn } from "./running-turns"; import { buildThreadTitle, isThreadTitled, @@ -65,16 +69,24 @@ function makeRefResolvers(workspaceId: number): RefResolvers { const logger = getLogger("api-server"); -const processedEvents = new Map(); +const DEDUP_TTL_SECONDS = 10 * 60; -function dedup(eventId: string): boolean { - const now = Date.now(); - for (const [id, ts] of processedEvents) { - if (now - ts > 300_000) processedEvents.delete(id); +/** + * Whether this event was already claimed, by this instance or another. If + * Redis is unreachable the event is processed: a duplicate turn beats a + * dropped one. + */ +async function isDuplicate(key: string): Promise { + try { + const claimed = await redis.set(`slack:event:${key}`, "1", { + nx: true, + ex: DEDUP_TTL_SECONDS, + }); + return claimed === null; + } catch (error) { + logger.warn("slack dedup unavailable", { error, key }); + return false; } - if (processedEvents.has(eventId)) return true; - processedEvents.set(eventId, now); - return false; } const slackEventSchema = z.object({ @@ -91,6 +103,13 @@ const slackEventSchema = z.object({ thread_ts: z.string().optional(), bot_id: z.string().optional(), tab: z.string().optional(), + // `tokens_revoked`: which tokens went. Only a revoked bot token matters. + tokens: z + .object({ + bot: z.array(z.string()).optional(), + oauth: z.array(z.string()).optional(), + }) + .optional(), assistant_thread: z .object({ channel_id: z.string(), @@ -233,7 +252,7 @@ export async function handleSlackEvent(c: Context) { return c.json({ ok: true }); } - if (body.event_id && dedup(body.event_id)) { + if (body.event_id && (await isDuplicate(body.event_id))) { return c.json({ ok: true }); } @@ -286,27 +305,13 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { if (event.type === "app_uninstalled" || event.type === "tokens_revoked") { const teamId = body.team_id; - if (teamId) { - await db - .delete(integration) - .where( - and( - eq(integration.name, "slack-agent"), - eq(integration.externalId, teamId), - ), - ); - await db - .update(pageSubscriber) - .set({ unsubscribedAt: new Date(), updatedAt: new Date() }) - .where( - and( - eq(pageSubscriber.channelType, "slack"), - isNull(pageSubscriber.unsubscribedAt), - sql`json_extract(${pageSubscriber.channelConfig}, '$.teamId') = ${teamId}`, - ), - ); - logger.info("slack integration cleaned up", { teamId }); - } + if (!teamId) return; + if (event.type === "tokens_revoked" && !event.tokens?.bot?.length) return; + const result = await uninstallSlackTeam({ + input: { teamId }, + job: `slack-${event.type}`, + }); + logger.info("slack integration cleaned up", { teamId, ...result }); return; } @@ -464,6 +469,9 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { if (!teamId || !channel || !threadTs) return; const wasRunning = abortTurn(channel, threadTs); + await broadcastStop(channel, threadTs).catch((error) => + logger.warn("slack failed to broadcast stop", { error, teamId }), + ); logger.info("slack turn stopped by user", { teamId, channel, @@ -515,8 +523,7 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { // A single mention arrives as BOTH an `app_mention` and a `message.*` event // (distinct event_ids, same message ts), so the event_id dedup above doesn't // catch the pair. Dedup on the message identity so we only respond once. - // Runs before the first `await` so concurrent deliveries can't both pass. - if (dedup(`msg:${event.channel}:${event.ts}`)) return; + if (await isDuplicate(`msg:${event.channel}:${event.ts}`)) return; const resolved = await resolveWorkspace(teamId); if (!resolved) { @@ -574,6 +581,24 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { const channel = event.channel; const slackUserId = event.user; if (!slackUserId) return; + if (!slackAgentAllowed(resolved.workspace)) { + if (await claimPlanNoticeWindow(teamId, slackUserId)) { + const message = planRequiredMessage(config); + await ( + isAgentThread + ? slack.chat.postMessage({ channel, thread_ts: threadTs, ...message }) + : slack.chat.postEphemeral({ + channel, + user: slackUserId, + thread_ts: event.thread_ts, + ...message, + }) + ).catch((error) => + logger.error("slack failed to send plan notice", { error, teamId }), + ); + } + return; + } const actor = await requireSlackMember({ workspace: resolved.workspace, teamId, @@ -615,7 +640,6 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { userId: event.user, isAgentThread, }); - if (!reply) return; if (turn.signal.aborted) { logger.info("slack turn stopped before it started", { @@ -631,9 +655,7 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { if (prefetchedThread) { thread = prefetchedThread; } else if (event.thread_ts) { - thread = ( - await fetchThread(slack, event.channel, event.thread_ts) - ).filter((msg) => msg.ts !== reply?.placeholderTs); + thread = await fetchThread(slack, event.channel, event.thread_ts); } else { thread = [{ user: event.user, text: event.text, ts: event.ts }]; } @@ -850,8 +872,7 @@ async function titleThread(args: { /** * Where the agent's output goes, in descending order of how much the user gets * to see while they wait: a streamed message that fills in as the model writes, - * with a task entry per tool call; a native "working" status on the thread; or - * a "Thinking..." message we post up front and overwrite. + * with a task entry per tool call; or a native "working" status on the thread. */ interface Reply { /** @@ -869,8 +890,6 @@ interface Reply { stopped(): Promise; /** Progress to report while the turn runs, when the surface can show it. */ progress?: AgentEvents; - /** Our own "Thinking..." message, to keep it out of the agent's context. */ - placeholderTs?: string; /** Runs once the turn is over, whether it succeeded or not. */ finish?: () => Promise; } @@ -1083,8 +1102,8 @@ async function createReply(args: { teamId: string; userId: string | undefined; isAgentThread: boolean; -}): Promise { - const { slack, channel, threadTs, teamId, userId, isAgentThread } = args; +}): Promise { + const { slack, channel, threadTs, teamId, userId } = args; const session = await acknowledgeWithSession( slack, @@ -1094,9 +1113,9 @@ async function createReply(args: { userId, ); - // The status is the loading indicator, not the delivery: a streamed turn in - // the agent pane still needs it when agent sessions aren't available. - if (!session && isAgentThread) { + // The status is the loading indicator, not the delivery: a streamed turn + // still needs it when agent sessions aren't available. + if (!session) { // Best-effort: a missing status only loses the loading indicator. await setAssistantStatus(slack, channel, threadTs, "is thinking...").catch( (err: unknown) => @@ -1120,9 +1139,7 @@ async function createReply(args: { }); } - if (session) return session; - if (isAgentThread) return acknowledgeInAgentThread(slack, channel, threadTs); - return acknowledgeInChannel(slack, channel, threadTs, teamId); + return session ?? acknowledgeWithStatus(slack, channel, threadTs); } function postInThread( @@ -1201,7 +1218,7 @@ async function acknowledgeWithSession( // The status was already set by `createReply`; Slack clears it as soon as the // app posts in the thread, so there is no matching "clear" call. -function acknowledgeInAgentThread( +function acknowledgeWithStatus( slack: WebClient, channel: string, threadTs: string, @@ -1216,79 +1233,6 @@ function acknowledgeInAgentThread( }; } -async function acknowledgeInChannel( - slack: WebClient, - channel: string, - threadTs: string, - teamId: string, -): Promise { - let thinkingTs: string | undefined; - try { - const thinkingMsg = await slack.chat.postMessage({ - channel, - thread_ts: threadTs, - text: ":hourglass_flowing_sand: Thinking...", - }); - thinkingTs = thinkingMsg.ts; - } catch (err) { - if (isSlackPlatformError(err, "cannot_reply_to_message")) { - logger.warn("slack cannot reply to message, falling back to top-level", { - channel, - teamId, - threadTs, - }); - try { - const fallbackMsg = await slack.chat.postMessage({ - channel, - text: ":hourglass_flowing_sand: Thinking...", - }); - thinkingTs = fallbackMsg.ts; - } catch (fallbackErr) { - logger.error("slack failed to post fallback thinking message", { - error: fallbackErr, - channel, - teamId, - }); - return; - } - } else { - logger.error("slack failed to post thinking message", { - error: err, - channel, - teamId, - threadTs, - }); - return; - } - } - - if (!thinkingTs) { - logger.error("slack thinking message returned no ts", { channel, teamId }); - return; - } - - const ts = thinkingTs; - const postNew = postInThread(slack, channel, threadTs); - let placeholderUsed = false; - // The first message takes over "Thinking..."; any after it (a second - // confirmation card) is a message of its own, or it would overwrite the first. - const send: Reply["send"] = async ({ text, blocks }) => { - if (placeholderUsed) return postNew({ text, blocks }); - placeholderUsed = true; - await slack.chat.update({ channel, ts, text, blocks }); - return ts; - }; - return { - placeholderTs: ts, - send, - answer: (text) => send(buildAnswerMessage(text)), - // Overwrites "Thinking...", which would otherwise stand forever. - stopped: async () => { - await send({ text: STOPPED_NOTICE }); - }, - }; -} - async function handleConfirmation( slack: WebClient, reply: Reply, @@ -1316,10 +1260,8 @@ async function handleConfirmation( // findByThread + replace isn't atomic on its own — two concurrent // events on the same thread could both see `existing` and race on - // replace. Atomicity here relies on the `dedup` map at the top of this - // file suppressing duplicate event_ids, plus Slack's own per-thread - // event throttling. Cross-process dedup is *not* covered; see note in - // processedEvents. + // replace. The Redis event dedup at the top of this file suppresses + // duplicate deliveries across instances, and Slack throttles per thread. const existing = await findByThread(threadTs, payload); if (existing) { await replace(existing.id, payload); diff --git a/apps/server/src/routes/slack/interactions.test.ts b/apps/server/src/routes/slack/interactions.test.ts index 75c0e503..09b9109c 100644 --- a/apps/server/src/routes/slack/interactions.test.ts +++ b/apps/server/src/routes/slack/interactions.test.ts @@ -48,7 +48,10 @@ function configureSlackDoubles() { }); slackTestState.resolveWorkspace = (teamId: string) => teamId === "T_KNOWN" - ? Promise.resolve({ botToken: "xoxb-fallback", workspace: { id: 1 } }) + ? Promise.resolve({ + botToken: "xoxb-fallback", + workspace: { id: 1, limits: { "slack-agent": true } }, + }) : Promise.resolve(null); } @@ -249,7 +252,10 @@ describe("handleSlackInteraction (dispatch)", () => { slackTestState.resolveWorkspace = (teamId: string) => { resolveCalls++; return teamId === "T_KNOWN" - ? Promise.resolve({ botToken: "xoxb-fresh", workspace: { id: 1 } }) + ? Promise.resolve({ + botToken: "xoxb-fresh", + workspace: { id: 1, limits: { "slack-agent": true } }, + }) : Promise.resolve(null); }; @@ -273,7 +279,10 @@ describe("handleSlackInteraction (dispatch)", () => { // A reinstall can point the team at a different workspace than the one the // card was drafted for; executing it there would hit the wrong status page. slackTestState.resolveWorkspace = () => - Promise.resolve({ botToken: "xoxb-other", workspace: { id: 2 } }); + Promise.resolve({ + botToken: "xoxb-other", + workspace: { id: 2, limits: { "slack-agent": true } }, + }); await signAndPost(app, { type: "block_actions", @@ -475,7 +484,7 @@ describe("registry-runner execution paths", () => { slackTestState.resolveWorkspace = () => Promise.resolve({ botToken: "xoxb-fallback", - workspace: { id: workspace.id }, + workspace: { id: workspace.id, limits: { "slack-agent": true } }, }); slackTestState.usersInfoImpl = () => Promise.resolve({ diff --git a/apps/server/src/routes/slack/interactions.ts b/apps/server/src/routes/slack/interactions.ts index 1398a87e..62bb1d2d 100644 --- a/apps/server/src/routes/slack/interactions.ts +++ b/apps/server/src/routes/slack/interactions.ts @@ -17,8 +17,10 @@ import { renderToolResult } from "./presenters"; import { executeRegistryAction, getRegistryTool } from "./registry-runner"; import { linkAccountUrl, + planRequiredMessage, requireSlackMember, type SlackActor, + slackAgentAllowed, } from "./require-slack-member"; import { toServiceCtx } from "./service-adapter"; import { resolveWorkspace } from "./workspace-resolver"; @@ -126,6 +128,14 @@ async function processInteraction( // exempt: the initiator can always dismiss their own draft. let actor: SlackActor | null = null; if (parsed.kind !== "cancel") { + if (!slackAgentAllowed(resolved.workspace)) { + await slack.chat.postEphemeral({ + channel: channelId, + user: userId, + ...planRequiredMessage(config), + }); + return; + } actor = await requireSlackMember({ workspace: resolved.workspace, teamId: workspaceTeamId, diff --git a/apps/server/src/routes/slack/require-slack-member.ts b/apps/server/src/routes/slack/require-slack-member.ts index a42ccf7f..5292527b 100644 --- a/apps/server/src/routes/slack/require-slack-member.ts +++ b/apps/server/src/routes/slack/require-slack-member.ts @@ -61,6 +61,19 @@ export async function claimLinkCardWindow( return claimed !== null; } +/** Claims the per-user plan-notice window, kept apart from the link card's. */ +export async function claimPlanNoticeWindow( + teamId: string, + slackUserId: string, +): Promise { + const claimed = await redis.set( + `slack:plannotice:${teamId}:${slackUserId}`, + "1", + { nx: true, ex: LINK_CARD_WINDOW_SECONDS }, + ); + return claimed !== null; +} + /** Frees the window again when the card never made it out. */ export async function releaseLinkCardWindow( teamId: string, @@ -68,3 +81,13 @@ export async function releaseLinkCardWindow( ): Promise { await redis.del(linkCardKey(teamId, slackUserId)); } + +export function slackAgentAllowed(workspace: Workspace): boolean { + return workspace.limits["slack-agent"] === true; +} + +export function planRequiredMessage(config: SlackConfig): { text: string } { + return { + text: `openstatus in Slack isn't included in this workspace's plan. <${config.dashboardUrl}/settings/billing|Upgrade> to use it.`, + }; +} diff --git a/apps/server/src/routes/slack/running-turns.ts b/apps/server/src/routes/slack/running-turns.ts index added666..71a0b450 100644 --- a/apps/server/src/routes/slack/running-turns.ts +++ b/apps/server/src/routes/slack/running-turns.ts @@ -1,27 +1,54 @@ +import { redis } from "@/libs/clients"; + /** * Turns currently being processed, so Slack's stop button can cancel them. * * Keyed by thread, which is what the user's stop acts on. A thread can hold - * more than one: the dedup in `handler.ts` is per message, so two messages - * sent in quick succession overlap, and stop means "stop this thread". - * - * Process-local: with more than one server instance the stop event can land - * where the turn isn't running, and that instance simply finds nothing to - * abort. The handler clears the session status either way, so the user always - * gets out of the loading state. + * more than one: two messages sent in quick succession overlap, and stop + * means "stop this thread". A stop is also written to Redis, and every + * running turn polls for it, so it reaches the instance running the turn. */ const turns = new Map>(); +const pollers = new Map>(); + +const STOP_POLL_MS = 1_000; +const STOP_TTL_SECONDS = 10 * 60; function key(channel: string, threadTs: string): string { return `${channel}:${threadTs}`; } +function stopKey(channel: string, threadTs: string): string { + return `slack:stop:${key(channel, threadTs)}`; +} + export function startTurn(channel: string, threadTs: string): AbortController { const controller = new AbortController(); const id = key(channel, threadTs); const running = turns.get(id); if (running) running.add(controller); else turns.set(id, new Set([controller])); + + // Stops are ordered by a Redis counter, not host clocks, so only a stop + // counted after this turn started aborts it. + const readStopSeq = async () => + Number((await redis.get(stopKey(channel, threadTs))) ?? 0); + let baseline: number | undefined; + readStopSeq() + .then((seq) => { + baseline ??= seq; + }) + .catch(() => {}); + const poller = setInterval(async () => { + try { + const seq = await readStopSeq(); + if (baseline === undefined) baseline = seq; + else if (seq > baseline) controller.abort(); + } catch { + // A Redis hiccup only delays a cross-instance stop until the next poll. + } + }, STOP_POLL_MS); + pollers.set(controller, poller); return controller; } @@ -30,6 +57,8 @@ export function endTurn( threadTs: string, controller: AbortController, ): void { + clearInterval(pollers.get(controller)); + pollers.delete(controller); const id = key(channel, threadTs); const running = turns.get(id); if (!running) return; @@ -44,3 +73,13 @@ export function abortTurn(channel: string, threadTs: string): boolean { for (const controller of running) controller.abort(); return true; } + +/** Tells turns of this thread running on other instances to stop. */ +export async function broadcastStop( + channel: string, + threadTs: string, +): Promise { + const k = stopKey(channel, threadTs); + await redis.incr(k); + await redis.expire(k, STOP_TTL_SECONDS); +} diff --git a/packages/services/src/integration/__tests__/uninstall-slack-agent.test.ts b/packages/services/src/integration/__tests__/uninstall-slack-agent.test.ts new file mode 100644 index 00000000..862c87f9 --- /dev/null +++ b/packages/services/src/integration/__tests__/uninstall-slack-agent.test.ts @@ -0,0 +1,130 @@ +import { and, eq } from "@openstatus/db"; +import { + integration, + pageSubscriber, + slackUser, +} from "@openstatus/db/src/schema"; +import { createPage, createSlackUser } from "@openstatus/db/src/test/factories"; +import { expect } from "@std/expect"; +import { describe, test } from "@std/testing/bdd"; + +import { + createWorkspaceFixture, + expectAuditRow, + makeSlackCtx, + withTestTransaction, +} from "../../../test/helpers"; +import { ForbiddenError } from "../../errors"; +import { createSlackSubscriber } from "../../page-subscriber/slack"; +import { installSlackAgent } from "../install-slack-agent"; +import { uninstallSlackTeam } from "../uninstall-slack-agent"; + +const input = (teamId: string) => ({ + externalId: teamId, + credential: { botToken: "xoxb-secret", botUserId: "B_BOT" }, + data: { + teamId, + teamName: "Fixture Team", + appId: "A_APP", + scopes: "chat:write", + installedBy: "U_INSTALLER", + }, +}); + +describe("installSlackAgent plan check", () => { + test("a plan without the Slack agent is refused", async () => { + const free = await createWorkspaceFixture("free"); + const ctx = makeSlackCtx(free.workspace, { + teamId: "T_FREE", + slackUserId: "U1", + userId: free.userId, + }); + await expect( + installSlackAgent({ ctx, input: input("T_FREE") }), + ).rejects.toThrow(ForbiddenError); + }); +}); + +describe("uninstallSlackTeam", () => { + test("removes the integration, links and subscriptions, all audited", async () => { + const teamId = `T_UNINSTALL_${crypto.randomUUID()}`; + const team = await createWorkspaceFixture("team"); + const page = await createPage(team.workspace.id); + + await withTestTransaction(async (tx) => { + const ctx = { + ...makeSlackCtx(team.workspace, { + teamId, + slackUserId: "U1", + userId: team.userId, + }), + db: tx, + }; + const installed = await installSlackAgent({ ctx, input: input(teamId) }); + const link = await createSlackUser( + team.workspace.id, + team.userId, + { slackTeamId: teamId }, + tx, + ); + const sub = await createSlackSubscriber({ + input: { pageId: page.id, teamId, channelId: "C_UNINSTALL" }, + db: tx, + }); + + const result = await uninstallSlackTeam({ + input: { teamId }, + job: "slack-app_uninstalled", + db: tx, + }); + expect(result).toEqual({ workspaces: 1, unsubscribed: 1 }); + + expect( + await tx + .select() + .from(integration) + .where(eq(integration.id, installed.id)) + .get(), + ).toBeUndefined(); + expect( + await tx + .select() + .from(slackUser) + .where(and(eq(slackUser.id, link.id))) + .get(), + ).toBeUndefined(); + const subscriber = await tx + .select() + .from(pageSubscriber) + .where(eq(pageSubscriber.id, sub.id)) + .get(); + expect(subscriber).toBeDefined(); + expect(subscriber?.unsubscribedAt).not.toBeNull(); + + await expectAuditRow({ + workspaceId: team.workspace.id, + action: "integration.delete", + entityType: "integration", + entityId: installed.id, + actorType: "system", + db: tx, + }); + await expectAuditRow({ + workspaceId: team.workspace.id, + action: "slack_user.delete", + entityType: "slack_user", + entityId: link.id, + actorType: "system", + db: tx, + }); + await expectAuditRow({ + workspaceId: team.workspace.id, + action: "page_subscriber.update", + entityType: "page_subscriber", + entityId: sub.id, + actorType: "system", + db: tx, + }); + }); + }); +}); diff --git a/packages/services/src/integration/index.ts b/packages/services/src/integration/index.ts index f6819105..893bf790 100644 --- a/packages/services/src/integration/index.ts +++ b/packages/services/src/integration/index.ts @@ -1,9 +1,14 @@ export { deleteIntegration } from "./delete"; export { installSlackAgent } from "./install-slack-agent"; +export { + uninstallSlackAgent, + uninstallSlackTeam, +} from "./uninstall-slack-agent"; export { listIntegrations, type IntegrationSummary } from "./list"; export { DeleteIntegrationInput, InstallSlackAgentInputSchema, type InstallSlackAgentInput, ListIntegrationsInput, + UninstallSlackTeamInput, } from "./schemas"; diff --git a/packages/services/src/integration/install-slack-agent.ts b/packages/services/src/integration/install-slack-agent.ts index 62961bfd..0f8a2611 100644 --- a/packages/services/src/integration/install-slack-agent.ts +++ b/packages/services/src/integration/install-slack-agent.ts @@ -4,7 +4,7 @@ import { integration } from "@openstatus/db/src/schema"; import { emitAudit } from "../audit"; import { requireScope } from "../auth"; import { type ServiceContext, withTransaction } from "../context"; -import { InternalServiceError, NotFoundError } from "../errors"; +import { ForbiddenError, InternalServiceError, NotFoundError } from "../errors"; import { type InstallSlackAgentInput, InstallSlackAgentInputSchema, @@ -27,6 +27,9 @@ export async function installSlackAgent(args: { const { ctx } = args; requireScope(ctx, "write"); const input = InstallSlackAgentInputSchema.parse(args.input); + if (!ctx.workspace.limits["slack-agent"]) { + throw new ForbiddenError("The Slack agent is not included in your plan"); + } return withTransaction(ctx, async (tx) => { // The DB has no `UNIQUE (workspace_id, name)` constraint, so a race diff --git a/packages/services/src/integration/schemas.ts b/packages/services/src/integration/schemas.ts index 5bc81561..a0716606 100644 --- a/packages/services/src/integration/schemas.ts +++ b/packages/services/src/integration/schemas.ts @@ -23,3 +23,6 @@ export const InstallSlackAgentInputSchema = z.object({ export type InstallSlackAgentInput = z.infer< typeof InstallSlackAgentInputSchema >; + +export const UninstallSlackTeamInput = z.object({ teamId: z.string().min(1) }); +export type UninstallSlackTeamInput = z.infer; diff --git a/packages/services/src/integration/uninstall-slack-agent.ts b/packages/services/src/integration/uninstall-slack-agent.ts new file mode 100644 index 00000000..32fdfe16 --- /dev/null +++ b/packages/services/src/integration/uninstall-slack-agent.ts @@ -0,0 +1,100 @@ +import { and, db as defaultDb, eq, inArray } from "@openstatus/db"; +import { integration, workspace } from "@openstatus/db/src/schema"; + +import { emitAudit } from "../audit"; +import { requireScope } from "../auth"; +import { type DB, type ServiceContext, withTransaction } from "../context"; +import { parseWorkspaceForContext } from "../page-subscriber/internal"; +import { removeSlackTeamSubscribers } from "../page-subscriber/slack"; +import { deleteSlackUserMappings } from "../slack-user/internal"; +import { UninstallSlackTeamInput } from "./schemas"; +import { snapshotIntegration } from "./snapshot"; + +/** + * Removes the workspace's link to a Slack team after Slack reports the app + * uninstalled or its bot token revoked: the integration row and the team's + * Slack account links go, each with its audit row. + */ +export async function uninstallSlackAgent(args: { + ctx: ServiceContext; + input: UninstallSlackTeamInput; +}): Promise { + const { ctx } = args; + requireScope(ctx, "write"); + const input = UninstallSlackTeamInput.parse(args.input); + + await withTransaction(ctx, async (tx) => { + const removed = await tx + .delete(integration) + .where( + and( + eq(integration.name, "slack-agent"), + eq(integration.workspaceId, ctx.workspace.id), + eq(integration.externalId, input.teamId), + ), + ) + .returning(); + for (const row of removed) { + await emitAudit(tx, ctx, { + action: "integration.delete", + entityType: "integration", + entityId: row.id, + before: await snapshotIntegration(row), + }); + } + await deleteSlackUserMappings({ + tx, + ctx, + where: { slackTeamId: input.teamId }, + }); + // TODO(incident/15-slack-channel): unbind every incident bound to this team. + }); +} + +/** + * Every workspace connected to the Slack team is uninstalled, and every + * channel of that team stops receiving status-page updates, whichever + * workspace owns the page. + */ +export async function uninstallSlackTeam(args: { + input: UninstallSlackTeamInput; + job: string; + db?: DB; +}): Promise<{ workspaces: number; unsubscribed: number }> { + const input = UninstallSlackTeamInput.parse(args.input); + const readDb = args.db ?? defaultDb; + const rows = await readDb + .select() + .from(workspace) + .where( + inArray( + workspace.id, + readDb + .select({ id: integration.workspaceId }) + .from(integration) + .where( + and( + eq(integration.name, "slack-agent"), + eq(integration.externalId, input.teamId), + ), + ), + ), + ) + .all(); + for (const row of rows) { + await uninstallSlackAgent({ + ctx: { + workspace: parseWorkspaceForContext(row), + actor: { type: "system", job: args.job }, + db: args.db, + }, + input, + }); + } + const unsubscribed = await removeSlackTeamSubscribers({ + input, + job: args.job, + db: args.db, + }); + return { workspaces: rows.length, unsubscribed }; +} diff --git a/packages/services/src/page-subscriber/slack.ts b/packages/services/src/page-subscriber/slack.ts index 5cf39b90..ba9057b9 100644 --- a/packages/services/src/page-subscriber/slack.ts +++ b/packages/services/src/page-subscriber/slack.ts @@ -226,6 +226,62 @@ export async function removeSlackSubscriber(args: { }); } +/** + * Unsubscribes every channel of a Slack team, on any workspace's page, when + * the app is removed from that team. Audited per subscriber as `system`. + */ +// Called by `uninstallSlackTeam` for a team, not a workspace: each row is +// audited in its page's workspace, and a system actor has no scope to check. +// oxlint-disable-next-line openstatus/services-mutation-guards +export async function removeSlackTeamSubscribers(args: { + input: { teamId: string }; + job: string; + db?: DB; +}): Promise { + const { teamId } = args.input; + return withTransaction({ db: args.db } as ServiceContext, async (tx) => { + const rows = await tx.query.pageSubscriber.findMany({ + where: and( + eq(pageSubscriber.channelType, "slack"), + isNull(pageSubscriber.unsubscribedAt), + sql`json_extract(${pageSubscriber.channelConfig}, '$.teamId') = ${teamId}`, + ), + with: { page: { with: { workspace: true } } }, + }); + for (const existing of rows) { + const updated = await tx + .update(pageSubscriber) + .set({ unsubscribedAt: new Date(), updatedAt: new Date() }) + .where(eq(pageSubscriber.id, existing.id)) + .returning() + .get(); + if (!existing.page?.workspace) continue; + const { page: _page, ...existingRow } = existing; + const { token: _bt, ...beforeSnap } = + selectPageSubscriberSchema.parse(existingRow); + const { token: _at, ...afterSnap } = selectPageSubscriberSchema.parse( + updated ?? existingRow, + ); + await emitAudit( + tx, + { + workspace: parseWorkspaceForContext(existing.page.workspace), + actor: { type: "system", job: args.job }, + db: tx, + }, + { + action: "page_subscriber.update", + entityType: "page_subscriber", + entityId: existing.id, + before: beforeSnap, + after: afterSnap, + }, + ); + } + return rows.length; + }); +} + export interface SlackSubscriptionSummary { id: number; pageId: number;