From 4f3d9049c8d9aa1449bc4bc60c3a6c787913d5d5 Mon Sep 17 00:00:00 2001 From: Thibault Le Ouay Date: Wed, 30 Sep 2026 15:27:17 +0200 Subject: [PATCH] workflows: stale-incident reminders (#2805) Co-authored-by: Claude Opus 5.5 --- apps/workflows/package.json | 1 + apps/workflows/src/cron/incident-reminders.ts | 41 +++ apps/workflows/src/cron/scheduler.ts | 9 +- .../src/incident/__tests__/reminders.test.ts | 298 ++++++++++++++++++ packages/services/src/incident/index.ts | 7 + packages/services/src/incident/reminders.ts | 210 ++++++++++++ pnpm-lock.yaml | 3 + 7 files changed, 568 insertions(+), 1 deletion(-) create mode 100644 apps/workflows/src/cron/incident-reminders.ts create mode 100644 packages/services/src/incident/__tests__/reminders.test.ts create mode 100644 packages/services/src/incident/reminders.ts 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..31ffbe08 --- /dev/null +++ b/apps/workflows/src/cron/incident-reminders.ts @@ -0,0 +1,41 @@ +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({ + // Default retries can hold the tick for ~30 min; a released claim retries + // on the next tick instead. + clientFor: (token) => new WebClient(token, { retryConfig: { retries: 2 } }), + claim: async (key, ttlSeconds) => + (await client.set(key, "1", { nx: true, ex: ttlSeconds })) !== null, + release: async (key) => { + await client.del(key); + }, + dashboardUrl: + env().NODE_ENV === "production" + ? "https://app.openstatus.dev" + : "http://localhost:3001", + }); + logger.info( + "incident reminders sent: {sent}, failed: {failed}, skipped: {skipped}", + { + sent: results.filter((r) => r.sent).length, + failed: results.filter((r) => r.target !== null && !r.sent).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..cf89978b --- /dev/null +++ b/packages/services/src/incident/__tests__/reminders.test.ts @@ -0,0 +1,298 @@ +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 { + claim: async (key: string) => { + if (keys.has(key)) return false; + keys.add(key); + return true; + }, + release: async (key: string) => { + keys.delete(key); + }, + }; +} + +async function connect( + tx: DB, + teamId: string, + scopes: readonly string[] = SLACK_BOT_SCOPES, +) { + await tx.insert(integration).values({ + name: "slack-agent", + workspaceId: workspace.id, + externalId: teamId, + credential: { botToken: "xoxb-test", botUserId: "UBOT" }, + data: { teamId, scopes: 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, + commanderId: userId, + 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 claims = memoryClaim(); + const run = () => + remindStaleIncidents({ + db: tx, + clientFor: () => client, + ...claims, + 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: reported once per window 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 claims = memoryClaim(); + const run = () => + remindStaleIncidents({ + db: tx, + clientFor: () => client, + ...claims, + dashboardUrl: "https://app.test", + }); + expect((await run()).find((r) => r.incidentId === row.id)).toEqual({ + incidentId: row.id, + target: null, + sent: false, + }); + expect( + (await run()).find((r) => r.incidentId === row.id), + ).toBeUndefined(); + expect(sent).toHaveLength(0); + }); + }); + + function failingClient(error: unknown) { + let calls = 0; + const client = { + chat: { + postMessage: async () => { + calls++; + if (calls === 1) throw error; + return { ok: true, ts: "1.1" }; + }, + }, + } as unknown as SlackIncidentClient; + return { client, calls: () => calls }; + } + + async function staleChannelIncident(tx: DB) { + const teamId = `T_REM_${crypto.randomUUID()}`; + await connect(tx, teamId); + return createIncident( + workspace.id, + { + declaredAt: new Date(Date.now() - 5 * HOUR), + slackTeamId: teamId, + slackChannelId: "C_FAIL", + }, + tx, + ); + } + + test("transient slack failure: not sent, retried on the next tick", async () => { + await withTestTransaction(async (tx) => { + const row = await staleChannelIncident(tx); + const { client, calls } = failingClient( + Object.assign(new Error("socket hang up"), { + code: "slack_webapi_request_error", + }), + ); + const claims = memoryClaim(); + const run = () => + remindStaleIncidents({ + db: tx, + clientFor: () => client, + ...claims, + dashboardUrl: "https://app.test", + }); + expect((await run()).find((r) => r.incidentId === row.id)?.sent).toBe( + false, + ); + expect((await run()).find((r) => r.incidentId === row.id)?.sent).toBe( + true, + ); + expect(calls()).toBe(2); + }); + }); + + test("slack platform error: not sent, window stays claimed", async () => { + await withTestTransaction(async (tx) => { + const row = await staleChannelIncident(tx); + const { client, calls } = failingClient( + Object.assign(new Error("channel_not_found"), { + code: "slack_webapi_platform_error", + }), + ); + const claims = memoryClaim(); + const run = () => + remindStaleIncidents({ + db: tx, + clientFor: () => client, + ...claims, + dashboardUrl: "https://app.test", + }); + expect((await run()).find((r) => r.incidentId === row.id)?.sent).toBe( + false, + ); + expect( + (await run()).find((r) => r.incidentId === row.id), + ).toBeUndefined(); + expect(calls()).toBe(1); + }); + }); + + test("install without chat:write: nothing claimed or sent", async () => { + await withTestTransaction(async (tx) => { + const teamId = `T_REM_${crypto.randomUUID()}`; + await connect( + tx, + teamId, + SLACK_BOT_SCOPES.filter((s) => s !== "chat:write"), + ); + const row = await createIncident( + workspace.id, + { + declaredAt: new Date(Date.now() - 5 * HOUR), + slackTeamId: teamId, + slackChannelId: "C_NOSCOPE", + }, + tx, + ); + const claimed: string[] = []; + const { client, sent } = fakeClient(); + const results = await remindStaleIncidents({ + db: tx, + clientFor: () => client, + claim: async (key) => { + claimed.push(key); + return true; + }, + release: async () => {}, + dashboardUrl: "https://app.test", + }); + expect(results.find((r) => r.incidentId === row.id)).toBeUndefined(); + expect( + claimed.some((k) => k.startsWith(`incident:reminder:${row.id}:`)), + ).toBe(false); + expect(sent).toHaveLength(0); + }); + }); +}); diff --git a/packages/services/src/incident/index.ts b/packages/services/src/incident/index.ts index c307352e..b5ccf26a 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..f97fc50b --- /dev/null +++ b/packages/services/src/incident/reminders.ts @@ -0,0 +1,210 @@ +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 { + type SlackConnection, + 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; + /** `false` when there was nobody to remind or Slack rejected the post. */ + sent: boolean; +}; + +/** `true` when this caller owns the key; used so each window fires once. */ +export type ClaimOnce = (key: string, ttlSeconds: number) => Promise; +export type ReleaseClaim = (key: string) => Promise; + +// The Slack client already retries network, 5xx and 429 errors; anything else +// (channel_not_found, not_in_channel, …) would fail again on the next tick. +function isPermanentSlackError(error: unknown): boolean { + return ( + typeof error === "object" && + error !== null && + "code" in error && + error.code === "slack_webapi_platform_error" + ); +} + +async function reminderConnection( + db: DB, + workspaceId: number, +): Promise { + const parsed = selectWorkspaceSchema.safeParse( + await db + .select() + .from(workspace) + .where(eq(workspace.id, workspaceId)) + .get(), + ); + if (!parsed.success) return null; + const ws = parsed.data; + if (!isFeatureEnabled(ws, INCIDENT_FEATURE)) return null; + if (!ws.limits["slack-agent"]) return null; + const ctx: ServiceContext = { + workspace: ws, + actor: { type: "system", job: "incident-reminders" }, + db, + }; + const connection = await getSlackConnection({ ctx }); + if (!connection || connection.missingScopes.includes("chat:write")) { + return null; + } + return connection; +} + +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; + release: ReleaseClaim; + 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[] = []; + const connections = new Map(); + for (const row of candidates) { + if (row.closedAt) continue; + if (!connections.has(row.workspaceId)) { + connections.set( + row.workspaceId, + await reminderConnection(db, row.workspaceId), + ); + } + const connection = connections.get(row.workspaceId); + if (!connection) 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; + + // Claimed before resolving the target so a skip is reported once per window. + const key = `incident:reminder:${row.id}:${last?.id ?? 0}:${window}`; + if (!(await args.claim(key, 14 * 24 * 60 * 60))) 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, sent: false }); + 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>`; + const sent = await args + .clientFor(connection.botToken) + .chat.postMessage({ + channel: + target.kind === "channel" ? target.channel : target.slackUserId, + text, + }) + .then( + () => true, + async (error) => { + console.warn("incident reminder failed", { + incidentId: row.id, + error, + }); + if (!isPermanentSlackError(error)) await args.release(key); + return false; + }, + ); + results.push({ incidentId: row.id, target, sent }); + } + 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 -- 2.51.2