diff --git a/apps/server/src/libs/test/doubles/slack-test-state.ts b/apps/server/src/libs/test/doubles/slack-test-state.ts index 3ded1791..2298f748 100644 --- a/apps/server/src/libs/test/doubles/slack-test-state.ts +++ b/apps/server/src/libs/test/doubles/slack-test-state.ts @@ -24,6 +24,8 @@ export interface SlackTestState { historyImpl: () => Promise; /** `users.info` result; the default has no email, so no mapping is created. */ usersInfoImpl: (args: Record) => Promise; + /** `reactions.get` result; the default message has no reactions. */ + reactionsGetImpl: (args: Record) => Promise; } const g = globalThis as Record; @@ -49,6 +51,7 @@ if (!g.__slackTestState) { messages: [{ user: "U1", text: "channel message", ts: "1.1" }], }), usersInfoImpl: () => Promise.resolve({ ok: true, user: { profile: {} } }), + reactionsGetImpl: () => Promise.resolve({ ok: true, message: {} }), } satisfies SlackTestState; } diff --git a/apps/server/src/libs/test/doubles/slack-web-api.mock.ts b/apps/server/src/libs/test/doubles/slack-web-api.mock.ts index d3a9078c..4e38b2d3 100644 --- a/apps/server/src/libs/test/doubles/slack-web-api.mock.ts +++ b/apps/server/src/libs/test/doubles/slack-web-api.mock.ts @@ -19,6 +19,23 @@ export class WebClient { s.calls.push({ method: "postEphemeral", args }); return Promise.resolve(); }, + getPermalink: (args: Record) => { + s.calls.push({ method: "chat.getPermalink", args }); + return Promise.resolve({ + ok: true, + permalink: `https://slack.test/archives/${args.channel}/p${args.message_ts}`, + }); + }, + }; + reactions = { + get: (args: Record) => { + s.calls.push({ method: "reactions.get", args }); + return s.reactionsGetImpl(args); + }, + add: (args: Record) => { + s.calls.push({ method: "reactions.add", args }); + return Promise.resolve({ ok: true }); + }, }; // Mirrors ChatStreamer: `ts` is undefined until the first append or stop. chatStream = (args: Record) => { @@ -56,6 +73,10 @@ export class WebClient { }; conversations = { replies: () => s.repliesImpl(), + join: (args: Record) => { + s.calls.push({ method: "conversations.join", args }); + return Promise.resolve({ ok: true }); + }, history: (args: Record) => { s.calls.push({ method: "conversations.history", args }); return s.historyImpl(); diff --git a/apps/server/src/routes/slack/handler.test.ts b/apps/server/src/routes/slack/handler.test.ts index b9bbbac4..4cc213aa 100644 --- a/apps/server/src/routes/slack/handler.test.ts +++ b/apps/server/src/routes/slack/handler.test.ts @@ -1,7 +1,7 @@ import crypto from "node:crypto"; import { db, eq } from "@openstatus/db"; -import { incident, integration } from "@openstatus/db/src/schema"; +import { incident, incidentEvent, 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"; @@ -2023,3 +2023,116 @@ describe("incident channel context", () => { } }); }); + +describe("incident channel events", () => { + const app = createTestApp(); + let incidentId: number; + let channelId: string; + + beforeEach(async () => { + resetSlackTestState(); + slackTestState.reactionsGetImpl = () => + Promise.resolve({ ok: true, message: {} }); + channelId = `C_PIN_${crypto.randomUUID()}`; + const [row] = await db + .insert(incident) + .values({ + workspaceId: 1, + title: "Pinned incident", + severity: "major", + declaredAt: new Date(), + startedAt: new Date(), + slackTeamId: "T_KNOWN", + slackChannelId: channelId, + }) + .returning(); + incidentId = row.id; + }); + + async function notes() { + return db + .select() + .from(incidentEvent) + .where(eq(incidentEvent.incidentId, incidentId)) + .all(); + } + + function pin(ts: string, user = "U1") { + return signAndPost(app, { + type: "event_callback", + team_id: "T_KNOWN", + event_id: `evt_pin_${crypto.randomUUID()}`, + event: { + type: "reaction_added", + user, + reaction: "pushpin", + item: { type: "message", channel: channelId, ts }, + }, + }); + } + + test("a 📌 copies the message onto the timeline and confirms with ✅", async () => { + slackTestState.historyImpl = () => + Promise.resolve({ + messages: [{ ts: "500.1", text: "Rolled back to v41", user: "U2" }], + }); + await pin("500.1"); + await waitForCall("reactions.add"); + const rows = await notes(); + expect(rows).toHaveLength(1); + expect(rows[0].message).toContain("Rolled back to v41"); + expect(rows[0].message).toContain("https://slack.test/archives/"); + }); + + test("a message already confirmed is not noted twice", async () => { + slackTestState.historyImpl = () => + Promise.resolve({ messages: [{ ts: "501.1", text: "Twice" }] }); + slackTestState.reactionsGetImpl = () => + Promise.resolve({ + ok: true, + message: { + reactions: [{ name: "white_check_mark", users: ["UBOT"] }], + }, + }); + await pin("501.1"); + await new Promise((r) => setTimeout(r, 150)); + expect(await notes()).toHaveLength(0); + }); + + test("a pinned thread reply is found through the thread", async () => { + slackTestState.historyImpl = () => + Promise.resolve({ messages: [{ ts: "400.0", text: "parent" }] }); + slackTestState.repliesImpl = () => + Promise.resolve({ + messages: [ + { ts: "400.0", text: "parent" }, + { ts: "400.5", text: "the reply" }, + ], + }); + await pin("400.5"); + await waitForCall("reactions.add"); + const rows = await notes(); + expect(rows[0]?.message).toContain("the reply"); + }); + + test("archiving the channel unbinds the incident", async () => { + await signAndPost(app, { + type: "event_callback", + team_id: "T_KNOWN", + event_id: `evt_archive_${crypto.randomUUID()}`, + event: { type: "channel_archive", channel: channelId, user: "U1" }, + }); + const deadline = Date.now() + 2000; + let bound: string | null = channelId; + while (bound && Date.now() < deadline) { + await new Promise((r) => setTimeout(r, 20)); + const row = await db + .select({ slackChannelId: incident.slackChannelId }) + .from(incident) + .where(eq(incident.id, incidentId)) + .get(); + bound = row?.slackChannelId ?? null; + } + expect(bound).toBeNull(); + }); +}); diff --git a/apps/server/src/routes/slack/handler.ts b/apps/server/src/routes/slack/handler.ts index 3cff0b6f..30760917 100644 --- a/apps/server/src/routes/slack/handler.ts +++ b/apps/server/src/routes/slack/handler.ts @@ -36,6 +36,7 @@ import { makeRefResolvers } from "./confirmation-card"; import { draftKey, findByThread, replace, store } from "./confirmation-store"; import type { PendingPayload } from "./confirmation-store"; import { publishHomeView, publishLinkAccountView } from "./home"; +import { handleChannelGone, handlePinReaction } from "./incident-events"; import { getComponentNames, getPageDashboardLink, @@ -101,6 +102,14 @@ const slackEventSchema = z.object({ thread_ts: z.string().optional(), bot_id: z.string().optional(), tab: z.string().optional(), + reaction: z.string().optional(), + item: z + .object({ + type: z.string(), + channel: z.string().optional(), + ts: z.string().optional(), + }) + .optional(), // `tokens_revoked`: which tokens went. Only a revoked bot token matters. tokens: z .object({ @@ -332,6 +341,39 @@ async function processEvent(body: SlackEvent, config: SlackConfig) { return; } + if (event.type === "reaction_added") { + const teamId = body.team_id; + const { item, user, reaction } = event; + if (!teamId || !user || !reaction || item?.type !== "message") return; + if (!item.channel || !item.ts) return; + const resolved = await resolveWorkspace(teamId); + if (!resolved) return; + await handlePinReaction({ + resolved, + config, + teamId, + slackUserId: user, + reaction, + channel: item.channel, + ts: item.ts, + }); + return; + } + + if (event.type === "channel_archive" || event.type === "channel_deleted") { + const teamId = body.team_id; + if (!teamId || !event.channel) return; + const resolved = await resolveWorkspace(teamId); + if (!resolved) return; + await handleChannelGone({ + resolved, + teamId, + channel: event.channel, + slackUserId: event.user, + }); + return; + } + // Both tabs of the app's DM arrive here. The agent experience has no // "thread started" event, so opening the Messages tab is where a first-time // user gets greeted; `greetOnce` makes the repeat opens harmless. diff --git a/apps/server/src/routes/slack/incident-events.ts b/apps/server/src/routes/slack/incident-events.ts new file mode 100644 index 00000000..bec10efa --- /dev/null +++ b/apps/server/src/routes/slack/incident-events.ts @@ -0,0 +1,169 @@ +import { getLogger } from "@logtape/logtape"; +import type { ServiceContext } from "@openstatus/services"; +import { + addIncidentNote, + getIncidentBySlackChannel, + unbindIncidentSlackChannel, +} from "@openstatus/services/incident"; +import { WebClient } from "@slack/web-api"; + +import { redis } from "@/libs/clients"; + +import { buildLinkAccountBlocks, LINK_ACCOUNT_TEXT } from "./blocks"; +import type { SlackConfig } from "./config"; +import { trackSlackIncident } from "./incident-analytics"; +import { + linkAccountUrl, + requireSlackMember, + slackAgentAllowed, +} from "./require-slack-member"; +import type { SlackWorkspace } from "./workspace-resolver"; + +const logger = getLogger(["api-server", "slack", "incident-events"]); + +const PIN = "pushpin"; +const DONE = "white_check_mark"; + +type SlackMessage = { ts?: string; text?: string; user?: string }; + +function system(resolved: SlackWorkspace): ServiceContext { + return { + workspace: resolved.workspace, + actor: { type: "system", job: "slack-incident-events" }, + }; +} + +async function findMessage( + slack: WebClient, + channel: string, + ts: string, +): Promise { + const history = await slack.conversations.history({ + channel, + latest: ts, + inclusive: true, + limit: 1, + }); + const top = (history.messages ?? []).find((m) => m.ts === ts); + if (top) return top; + // A thread reply is not in the channel history; ask the thread for it. + const replies = await slack.conversations.replies({ + channel, + ts, + latest: ts, + inclusive: true, + limit: 1, + }); + return (replies.messages ?? []).find((m) => m.ts === ts); +} + +async function alreadyNoted( + slack: WebClient, + channel: string, + ts: string, + botUserId: string, +): Promise { + const res = await slack.reactions.get({ channel, timestamp: ts }); + return (res.message?.reactions ?? []).some( + (r) => r.name === DONE && (r.users ?? []).includes(botUserId), + ); +} + +/** 📌 on a message in an incident channel copies it onto the timeline. */ +export async function handlePinReaction(args: { + resolved: SlackWorkspace; + config: SlackConfig; + teamId: string; + slackUserId: string; + reaction: string; + channel: string; + ts: string; +}): Promise { + const { resolved, config, teamId, slackUserId, channel, ts } = args; + if (args.reaction !== PIN) return; + if (!slackAgentAllowed(resolved.workspace)) return; + const bound = await getIncidentBySlackChannel({ + ctx: system(resolved), + input: { teamId, channelId: channel }, + }); + if (!bound || bound.closedAt) return; + + const slack = new WebClient(resolved.botToken); + const actor = await requireSlackMember({ + workspace: resolved.workspace, + teamId, + slackUserId, + slack, + }); + if (!actor) { + const once = await redis.set(`slack:pinlink:${channel}:${ts}`, "1", { + nx: true, + ex: 24 * 60 * 60, + }); + if (once === null) return; + const url = await linkAccountUrl(config, { + workspaceId: resolved.workspace.id, + teamId, + slackUserId, + }); + await slack.chat.postEphemeral({ + channel, + user: slackUserId, + thread_ts: ts, + text: LINK_ACCOUNT_TEXT, + blocks: buildLinkAccountBlocks(url), + }); + return; + } + + if (await alreadyNoted(slack, channel, ts, resolved.botUserId)) return; + const message = await findMessage(slack, channel, ts); + if (!message?.text) return; + const permalink = await slack.chat + .getPermalink({ channel, message_ts: ts }) + .then((res) => res.permalink) + .catch(() => undefined); + + const ctx: ServiceContext = { workspace: resolved.workspace, actor }; + await addIncidentNote({ + ctx, + input: { + id: bound.id, + message: permalink + ? `${message.text}\n\n[From Slack](${permalink})` + : message.text, + }, + }); + trackSlackIncident(ctx, "note", { via: "reaction" }); + await slack.reactions + .add({ channel, timestamp: ts, name: DONE }) + .catch((error) => logger.warn("slack failed to confirm pin", { error })); +} + +/** An archived or deleted channel no longer carries its incident. */ +export async function handleChannelGone(args: { + resolved: SlackWorkspace; + teamId: string; + channel: string; + slackUserId: string | undefined; +}): Promise { + const { resolved, teamId, channel } = args; + const bound = await getIncidentBySlackChannel({ + ctx: system(resolved), + input: { teamId, channelId: channel }, + }); + if (!bound) return; + const actor = args.slackUserId + ? await requireSlackMember({ + workspace: resolved.workspace, + teamId, + slackUserId: args.slackUserId, + slack: new WebClient(resolved.botToken), + }) + : null; + await unbindIncidentSlackChannel({ + ctx: actor ? { workspace: resolved.workspace, actor } : system(resolved), + input: { id: bound.id }, + }); + logger.info("slack incident channel unbound", { incidentId: bound.id }); +}