diff --git a/apps/workflows/package.json b/apps/workflows/package.json index dd8f9f09..46624fca 100644 --- a/apps/workflows/package.json +++ b/apps/workflows/package.json @@ -44,6 +44,7 @@ "@opentelemetry/sdk-logs": "catalog:", "@opentelemetry/semantic-conventions": "catalog:", "@sentry/deno": "catalog:", + "@slack/web-api": "catalog:", "@upstash/qstash": "catalog:", "drizzle-orm": "catalog:", "effect": "catalog:", diff --git a/apps/workflows/src/cron/incident-reminders.ts b/apps/workflows/src/cron/incident-reminders.ts new file mode 100644 index 00000000..5829de7a --- /dev/null +++ b/apps/workflows/src/cron/incident-reminders.ts @@ -0,0 +1,32 @@ +import { getLogger } from "@logtape/logtape"; +import { + type ReminderResult, + remindStaleIncidents, +} from "@openstatus/services/incident"; +import { Redis } from "@openstatus/upstash"; +import { WebClient } from "@slack/web-api"; + +import { env } from "../env"; + +const logger = getLogger(["workflow", "incident-reminders"]); + +let redis: ReturnType | undefined; + +export async function runIncidentRemindersTick(): Promise { + redis ??= Redis.fromEnv(); + const client = redis; + const results = await remindStaleIncidents({ + clientFor: (token) => new WebClient(token), + claim: async (key, ttlSeconds) => + (await client.set(key, "1", { nx: true, ex: ttlSeconds })) !== null, + dashboardUrl: + env().NODE_ENV === "production" + ? "https://app.openstatus.dev" + : "http://localhost:3001", + }); + logger.info("incident reminders sent: {sent}, skipped: {skipped}", { + sent: results.filter((r) => r.target !== null).length, + skipped: results.filter((r) => r.target === null).length, + }); + return results; +} diff --git a/apps/workflows/src/cron/scheduler.ts b/apps/workflows/src/cron/scheduler.ts index 41141ee4..ccda55f9 100644 --- a/apps/workflows/src/cron/scheduler.ts +++ b/apps/workflows/src/cron/scheduler.ts @@ -2,6 +2,7 @@ import { getLogger } from "@logtape/logtape"; import * as Sentry from "@sentry/deno"; import { Effect, Fiber, Schedule } from "effect"; +import { runIncidentRemindersTick } from "./incident-reminders"; import { runOAuthPruneTick } from "./oauth-prune"; import { handleOutboxDrainCron, handleOutboxRetentionCron } from "./outbox"; import { handleStatusDriftCron } from "./status-drift"; @@ -18,7 +19,8 @@ type ScheduledTask = { * Internal maintenance only, so it runs in-process rather than needing a * schedule added outside this repo. Every task is safe to run on both machines * at once: the outbox claim is atomic, drift repair is guarded by the same - * compare-and-swap as a live check, and retention deletes are idempotent. + * compare-and-swap as a live check, retention deletes are idempotent, and each + * incident reminder window is claimed in Redis before it is sent. */ export const SCHEDULED_TASKS: ScheduledTask[] = [ { @@ -41,6 +43,11 @@ export const SCHEDULED_TASKS: ScheduledTask[] = [ expression: "41 3 * * *", run: runOAuthPruneTick, }, + { + name: "incident-reminders", + expression: "*/10 * * * *", + run: runIncidentRemindersTick, + }, ]; type RunningTask = Fiber.Fiber; diff --git a/packages/services/src/incident/__tests__/reminders.test.ts b/packages/services/src/incident/__tests__/reminders.test.ts new file mode 100644 index 00000000..398452d0 --- /dev/null +++ b/packages/services/src/incident/__tests__/reminders.test.ts @@ -0,0 +1,162 @@ +import { integration } from "@openstatus/db/src/schema"; +import { + createIncident, + createIncidentEvent, + createSlackUser, +} from "@openstatus/db/src/test/factories"; +import { expect } from "@std/expect"; +import { beforeAll, describe, test } from "@std/testing/bdd"; + +import { + createWorkspaceFixture, + withTestTransaction, +} from "../../../test/helpers"; +import type { DB } from "../../context"; +import { SLACK_BOT_SCOPES } from "../../integration/slack-scopes"; +import type { Workspace } from "../../types"; +import { + reminderWindow, + remindStaleIncidents, + type SlackIncidentClient, +} from "../index"; + +const HOUR = 60 * 60 * 1000; + +let workspace: Workspace; +let userId: number; + +beforeAll(async () => { + const fixture = await createWorkspaceFixture("team"); + workspace = fixture.workspace; + userId = fixture.userId; +}); + +describe("reminderWindow", () => { + const now = new Date("2026-09-28T12:00:00Z"); + const ago = (h: number) => new Date(now.getTime() - h * HOUR); + + test("fresh, then windows that double", () => { + expect(reminderWindow("major", ago(3), now)).toBeNull(); + expect(reminderWindow("major", ago(4), now)).toBe(0); + expect(reminderWindow("major", ago(7.9), now)).toBe(0); + expect(reminderWindow("major", ago(8), now)).toBe(1); + expect(reminderWindow("major", ago(16), now)).toBe(2); + expect(reminderWindow("critical", ago(1), now)).toBe(0); + expect(reminderWindow("minor", ago(23), now)).toBeNull(); + }); +}); + +function fakeClient() { + const sent: { channel: string; text: string }[] = []; + const client = { + chat: { + postMessage: async (args: { channel: string; text: string }) => { + sent.push(args); + return { ok: true, ts: "1.1" }; + }, + }, + } as SlackIncidentClient; + return { client, sent }; +} + +function memoryClaim() { + const keys = new Set(); + return async (key: string) => { + if (keys.has(key)) return false; + keys.add(key); + return true; + }; +} + +async function connect(tx: DB, teamId: string) { + await tx.insert(integration).values({ + name: "slack-agent", + workspaceId: workspace.id, + externalId: teamId, + credential: { botToken: "xoxb-test", botUserId: "UBOT" }, + data: { teamId, scopes: SLACK_BOT_SCOPES.join(",") }, + }); +} + +describe("remindStaleIncidents", () => { + test("bound channel first, then the commander's DM, once per window", async () => { + await withTestTransaction(async (tx) => { + const teamId = `T_REM_${crypto.randomUUID()}`; + await connect(tx, teamId); + const stale = new Date(Date.now() - 5 * HOUR); + const inChannel = await createIncident( + workspace.id, + { + declaredAt: stale, + slackTeamId: teamId, + slackChannelId: "C_REM", + }, + tx, + ); + const byDm = await createIncident( + workspace.id, + { declaredAt: stale, commanderId: userId }, + tx, + ); + await createIncidentEvent(byDm.id, { createdAt: stale }, tx); + await createSlackUser( + workspace.id, + userId, + { slackTeamId: teamId, slackUserId: "U_COMMANDER" }, + tx, + ); + const fresh = await createIncident(workspace.id, {}, tx); + await createIncidentEvent(fresh.id, { createdAt: new Date() }, tx); + + const { client, sent } = fakeClient(); + const claim = memoryClaim(); + const run = () => + remindStaleIncidents({ + db: tx, + clientFor: () => client, + claim, + dashboardUrl: "https://app.test", + }); + + const results = await run(); + const mine = (id: number) => results.find((r) => r.incidentId === id); + expect(mine(inChannel.id)?.target).toEqual({ + kind: "channel", + channel: "C_REM", + }); + expect(mine(byDm.id)?.target).toEqual({ + kind: "dm", + slackUserId: "U_COMMANDER", + role: "commander", + }); + expect(mine(fresh.id)).toBeUndefined(); + expect(sent.map((m) => m.channel).sort()).toEqual( + ["C_REM", "U_COMMANDER"].sort(), + ); + + await run(); + expect(sent).toHaveLength(2); + }); + }); + + test("nobody linked: skipped and reported with no target", async () => { + await withTestTransaction(async (tx) => { + const teamId = `T_REM_${crypto.randomUUID()}`; + await connect(tx, teamId); + const row = await createIncident( + workspace.id, + { declaredAt: new Date(Date.now() - 30 * HOUR), severity: "minor" }, + tx, + ); + const { client, sent } = fakeClient(); + const results = await remindStaleIncidents({ + db: tx, + clientFor: () => client, + claim: memoryClaim(), + dashboardUrl: "https://app.test", + }); + expect(results.find((r) => r.incidentId === row.id)?.target).toBeNull(); + expect(sent).toHaveLength(0); + }); + }); +}); diff --git a/packages/services/src/incident/index.ts b/packages/services/src/incident/index.ts index de9726b0..c50fb1b0 100644 --- a/packages/services/src/incident/index.ts +++ b/packages/services/src/incident/index.ts @@ -10,6 +10,13 @@ export { } from "./link-status-report"; export { listIncidentEvents } from "./list-events"; export { clearIncidentCommander } from "./members"; +export { + type ClaimOnce, + reminderWindow, + remindStaleIncidents, + type ReminderResult, + STALE_AFTER, +} from "./reminders"; export { getIncident, getIncidentBySlackChannel, diff --git a/packages/services/src/incident/reminders.ts b/packages/services/src/incident/reminders.ts new file mode 100644 index 00000000..55c98503 --- /dev/null +++ b/packages/services/src/incident/reminders.ts @@ -0,0 +1,165 @@ +import { and, db as defaultDb, desc, eq, inArray } from "@openstatus/db"; +import { + type Incident, + type IncidentSeverity, + incident, + incidentEvent, + slackUser, + selectWorkspaceSchema, + workspace, +} from "@openstatus/db/src/schema"; + +import type { DB, ServiceContext } from "../context"; +import { isFeatureEnabled } from "../features"; +import { getSlackConnection } from "../integration/slack-connection"; +import { INCIDENT_FEATURE } from "./internal"; +import type { SlackClientFactory } from "./slack-flow"; + +const HOUR = 60 * 60 * 1000; + +export const STALE_AFTER: Record = { + critical: HOUR, + major: 4 * HOUR, + minor: 24 * HOUR, +}; + +/** + * Which reminder window an incident is in: 0 once it's stale, then 1, 2, … + * each twice as long as the last. `null` while it's still fresh. + */ +export function reminderWindow( + severity: IncidentSeverity, + lastActivity: Date, + now: Date, +): number | null { + const threshold = STALE_AFTER[severity]; + const elapsed = now.getTime() - lastActivity.getTime(); + if (elapsed < threshold) return null; + return Math.floor(Math.log2(elapsed / threshold)); +} + +export type ReminderTarget = + | { kind: "channel"; channel: string } + | { kind: "dm"; slackUserId: string; role: "commander" | "declarer" }; + +export type ReminderResult = { + incidentId: number; + target: ReminderTarget | null; +}; + +/** `true` when this caller owns the key; used so each window fires once. */ +export type ClaimOnce = (key: string, ttlSeconds: number) => Promise; + +async function dmTarget( + db: DB, + row: Incident, + teamId: string, +): Promise { + for (const [role, userId] of [ + ["commander", row.commanderId], + ["declarer", row.declaredBy], + ] as const) { + if (userId === null) continue; + const mapping = await db + .select({ slackUserId: slackUser.slackUserId }) + .from(slackUser) + .where( + and( + eq(slackUser.workspaceId, row.workspaceId), + eq(slackUser.slackTeamId, teamId), + eq(slackUser.userId, userId), + ), + ) + .get(); + if (mapping) return { kind: "dm", slackUserId: mapping.slackUserId, role }; + } + return null; +} + +/** + * Nudges open and mitigated incidents that went quiet: the bound channel, + * else a DM to the commander, else to the declarer. Safe to run on every + * machine at once: each window is claimed before anything is sent. + */ +export async function remindStaleIncidents(args: { + now?: Date; + clientFor: SlackClientFactory; + claim: ClaimOnce; + dashboardUrl: string; + db?: DB; +}): Promise { + const db = args.db ?? defaultDb; + const now = args.now ?? new Date(); + const candidates = await db + .select() + .from(incident) + .where(inArray(incident.status, ["open", "mitigated"])) + .all(); + + const results: ReminderResult[] = []; + for (const row of candidates) { + if (row.closedAt) continue; + const workspaceRow = await db + .select() + .from(workspace) + .where(eq(workspace.id, row.workspaceId)) + .get(); + const parsed = selectWorkspaceSchema.safeParse(workspaceRow); + if (!parsed.success) continue; + const ws = parsed.data; + if (!isFeatureEnabled(ws, INCIDENT_FEATURE)) continue; + if (!ws.limits["slack-agent"]) continue; + + const last = await db + .select({ id: incidentEvent.id, createdAt: incidentEvent.createdAt }) + .from(incidentEvent) + .where(eq(incidentEvent.incidentId, row.id)) + .orderBy(desc(incidentEvent.createdAt), desc(incidentEvent.id)) + .get(); + const lastActivity = last?.createdAt ?? row.declaredAt; + const window = reminderWindow(row.severity, lastActivity, now); + if (window === null) continue; + + const ctx: ServiceContext = { + workspace: ws, + actor: { type: "system", job: "incident-reminders" }, + db, + }; + const connection = await getSlackConnection({ ctx }); + if (!connection) continue; + + const target: ReminderTarget | null = + row.slackChannelId && row.slackTeamId === connection.teamId + ? { kind: "channel", channel: row.slackChannelId } + : await dmTarget(db, row, connection.teamId); + if (!target) { + console.warn("incident reminder skipped: nobody to remind", { + incidentId: row.id, + }); + results.push({ incidentId: row.id, target: null }); + continue; + } + + const key = `incident:reminder:${row.id}:${last?.id ?? 0}:${window}`; + if (!(await args.claim(key, 14 * 24 * 60 * 60))) continue; + + const quiet = Math.round((now.getTime() - lastActivity.getTime()) / HOUR); + const url = `${args.dashboardUrl}/incidents/${row.id}`; + const text = + target.kind === "channel" + ? `:hourglass: No update on *${row.title}* for ${quiet}h. Post a note, change its status, or resolve it. <${url}|Open in openstatus>` + : `:hourglass: You are the ${target.role} of *${row.title}* (${row.severity}), which has had no update for ${quiet}h. <${url}|Open in openstatus>`; + await args + .clientFor(connection.botToken) + .chat.postMessage({ + channel: + target.kind === "channel" ? target.channel : target.slackUserId, + text, + }) + .catch((error) => + console.warn("incident reminder failed", { incidentId: row.id, error }), + ); + results.push({ incidentId: row.id, target }); + } + return results; +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 30564427..c15e2595 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1890,6 +1890,9 @@ importers: '@sentry/deno': specifier: 'catalog:' version: 10.75.1 + '@slack/web-api': + specifier: 'catalog:' + version: 8.1.1 '@upstash/qstash': specifier: 'catalog:' version: 2.11.3