diff --git a/packages/subscriptions/src/channels/slack-blocks.test.ts b/packages/subscriptions/src/channels/slack-blocks.test.ts index 83dd0dcba..d33ece0ad 100644 --- a/packages/subscriptions/src/channels/slack-blocks.test.ts +++ b/packages/subscriptions/src/channels/slack-blocks.test.ts @@ -75,6 +75,21 @@ describe("buildRootMessage", () => { expect(contextText(maintenance)).not.toContain("Updated"); }); + test("maintenance root pins the scheduled start, not the update date", () => { + const root = buildRootMessage( + makeUpdate({ + status: "maintenance", + date: "2026-01-01T12:30:00.000Z", + startsAt: "2026-01-02T00:00:00.000Z", + endsAt: "2026-01-02T02:00:00.000Z", + }), + makeSub(), + ); + const text = JSON.stringify(root.attachments[0]?.blocks); + expect(text).toContain("Scheduled 2026-01-02T00:00:00.000Z"); + expect(text).not.toContain("2026-01-01T12:30:00.000Z"); + }); + test("uses the custom domain origin when present", () => { const root = buildRootMessage( makeUpdate(), diff --git a/packages/subscriptions/src/channels/slack-blocks.ts b/packages/subscriptions/src/channels/slack-blocks.ts index 0531fcce9..542c82aff 100644 --- a/packages/subscriptions/src/channels/slack-blocks.ts +++ b/packages/subscriptions/src/channels/slack-blocks.ts @@ -95,15 +95,19 @@ export function buildRootMessage( }); } - // Maintenance carries a scheduled window, not an update timestamp. - const dateLabel = - pageUpdate.status === "maintenance" ? "Scheduled" : "Updated"; + // Maintenance carries a scheduled window, not an update timestamp. The root + // is re-rendered on every update, so it must not pick up the update's date. + const isMaintenance = pageUpdate.status === "maintenance"; + const dateLabel = isMaintenance ? "Scheduled" : "Updated"; + const dateValue = isMaintenance + ? (pageUpdate.startsAt ?? pageUpdate.date) + : pageUpdate.date; blocks.push({ type: "context", elements: [ { type: "mrkdwn", - text: `${dateLabel} ${pageUpdate.date} · <${eventUrl(pageUpdate, subscription)}|View details> · Manage with \`/openstatus unsubscribe\``, + text: `${dateLabel} ${dateValue} · <${eventUrl(pageUpdate, subscription)}|View details> · Manage with \`/openstatus unsubscribe\``, }, ], }); diff --git a/packages/subscriptions/src/channels/slack-store.ts b/packages/subscriptions/src/channels/slack-store.ts index a635cf809..b7e9b1e40 100644 --- a/packages/subscriptions/src/channels/slack-store.ts +++ b/packages/subscriptions/src/channels/slack-store.ts @@ -11,29 +11,36 @@ export interface SlackThreadAnchor { pendingRootReply?: SlackReplyMessage; } +// The entity a thread belongs to. Report and maintenance ids come from +// different tables, so the kind keeps their anchors apart. +export interface SlackThreadEvent { + kind: "report" | "maintenance"; + id: number; +} + export interface SlackAnchorStore { getAnchor( - reportId: number, + event: SlackThreadEvent, subscriberId: number, ): Promise; setAnchor( - reportId: number, + event: SlackThreadEvent, subscriberId: number, anchor: SlackThreadAnchor, ): Promise; - clearAnchor(reportId: number, subscriberId: number): Promise; - // Atomically claim delivery of (report, subscriber, update). Returns true only + clearAnchor(event: SlackThreadEvent, subscriberId: number): Promise; + // Atomically claim delivery of (event, subscriber, update). Returns true only // for the caller that wins the claim; concurrent callers get false and must // skip. This is the dedupe reservation — a single atomic op, not a // read-then-write pair, so two dispatchers can't both post the same message. reserveDelivery( - reportId: number, + event: SlackThreadEvent, subscriberId: number, updateId: number, ): Promise; // Release a reservation whose post failed, so the delivery can be retried. releaseDelivery( - reportId: number, + event: SlackThreadEvent, subscriberId: number, updateId: number, ): Promise; @@ -41,16 +48,16 @@ export interface SlackAnchorStore { const TTL_SECONDS = 90 * 24 * 60 * 60; -function anchorKey(reportId: number, subscriberId: number): string { - return `slack:report:${reportId}:sub:${subscriberId}`; +function anchorKey(event: SlackThreadEvent, subscriberId: number): string { + return `slack:${event.kind}:${event.id}:sub:${subscriberId}`; } function deliveredKey( - reportId: number, + event: SlackThreadEvent, subscriberId: number, updateId: number, ): string { - return `${anchorKey(reportId, subscriberId)}:update:${updateId}`; + return `${anchorKey(event, subscriberId)}:update:${updateId}`; } let redisClient: Redis | null = null; @@ -64,30 +71,30 @@ function getRedis(): Redis { export function createRedisAnchorStore(): SlackAnchorStore { return { - async getAnchor(reportId, subscriberId) { + async getAnchor(event, subscriberId) { const raw = await getRedis().get( - anchorKey(reportId, subscriberId), + anchorKey(event, subscriberId), ); return raw ?? null; }, - async setAnchor(reportId, subscriberId, anchor) { - await getRedis().set(anchorKey(reportId, subscriberId), anchor, { + async setAnchor(event, subscriberId, anchor) { + await getRedis().set(anchorKey(event, subscriberId), anchor, { ex: TTL_SECONDS, }); }, - async clearAnchor(reportId, subscriberId) { - await getRedis().del(anchorKey(reportId, subscriberId)); + async clearAnchor(event, subscriberId) { + await getRedis().del(anchorKey(event, subscriberId)); }, - async reserveDelivery(reportId, subscriberId, updateId) { + async reserveDelivery(event, subscriberId, updateId) { const res = await getRedis().set( - deliveredKey(reportId, subscriberId, updateId), + deliveredKey(event, subscriberId, updateId), 1, { ex: TTL_SECONDS, nx: true }, ); return res === "OK"; }, - async releaseDelivery(reportId, subscriberId, updateId) { - await getRedis().del(deliveredKey(reportId, subscriberId, updateId)); + async releaseDelivery(event, subscriberId, updateId) { + await getRedis().del(deliveredKey(event, subscriberId, updateId)); }, }; } @@ -96,23 +103,23 @@ export function createMemoryAnchorStore(): SlackAnchorStore { const anchors = new Map(); const delivered = new Set(); return { - async getAnchor(reportId, subscriberId) { - return anchors.get(anchorKey(reportId, subscriberId)) ?? null; + async getAnchor(event, subscriberId) { + return anchors.get(anchorKey(event, subscriberId)) ?? null; }, - async setAnchor(reportId, subscriberId, anchor) { - anchors.set(anchorKey(reportId, subscriberId), anchor); + async setAnchor(event, subscriberId, anchor) { + anchors.set(anchorKey(event, subscriberId), anchor); }, - async clearAnchor(reportId, subscriberId) { - anchors.delete(anchorKey(reportId, subscriberId)); + async clearAnchor(event, subscriberId) { + anchors.delete(anchorKey(event, subscriberId)); }, - async reserveDelivery(reportId, subscriberId, updateId) { - const key = deliveredKey(reportId, subscriberId, updateId); + async reserveDelivery(event, subscriberId, updateId) { + const key = deliveredKey(event, subscriberId, updateId); if (delivered.has(key)) return false; delivered.add(key); return true; }, - async releaseDelivery(reportId, subscriberId, updateId) { - delivered.delete(deliveredKey(reportId, subscriberId, updateId)); + async releaseDelivery(event, subscriberId, updateId) { + delivered.delete(deliveredKey(event, subscriberId, updateId)); }, }; } diff --git a/packages/subscriptions/src/channels/slack.test.ts b/packages/subscriptions/src/channels/slack.test.ts index 8dabbe03b..48903b3c3 100644 --- a/packages/subscriptions/src/channels/slack.test.ts +++ b/packages/subscriptions/src/channels/slack.test.ts @@ -78,6 +78,8 @@ function makeUpdate(over: Partial = {}): PageUpdate { const noopUnsub = async () => {}; const token = async () => "xoxb-test"; +const REPORT = { kind: "report", id: 10 } as const; +const MAINTENANCE = { kind: "maintenance", id: 10 } as const; describe("createSlackChannel", () => { test("first update opens the thread (root post, no update, anchor set)", async () => { @@ -95,7 +97,7 @@ describe("createSlackChannel", () => { expect(calls.filter((c) => c.method === "post").length).toBe(1); expect(calls[0]?.thread_ts).toBeUndefined(); expect(calls.some((c) => c.method === "update")).toBe(false); - expect(await store.getAnchor(10, 1)).not.toBeNull(); + expect(await store.getAnchor(REPORT, 1)).not.toBeNull(); }); test("subsequent update backfills the first update then replies in thread and re-renders root", async () => { @@ -267,7 +269,7 @@ describe("createSlackChannel", () => { expect(calls.filter((c) => c.method === "post").length).toBe(1); }); - test("maintenance posts once with no thread and no anchor", async () => { + test("maintenance threads like a report: root, then backfill + reply + re-render", async () => { const { client, calls } = makeClient(); const store = createMemoryAnchorStore(); const channel = createSlackChannel({ @@ -276,15 +278,57 @@ describe("createSlackChannel", () => { getBotToken: token, softUnsubscribe: noopUnsub, }); + const maintenance = { + status: "maintenance" as const, + startsAt: "2026-01-02T00:00:00.000Z", + endsAt: "2026-01-02T02:00:00.000Z", + }; await channel.sendNotifications( [makeSub()], - makeUpdate({ status: "maintenance", updateId: undefined }), + makeUpdate({ ...maintenance, updateId: 200, message: "scheduled" }), ); - expect(calls.filter((c) => c.method === "post").length).toBe(1); + expect(calls[0]?.thread_ts).toBeUndefined(); + expect(await store.getAnchor(MAINTENANCE, 1)).not.toBeNull(); + + await channel.sendNotifications( + [makeSub()], + makeUpdate({ ...maintenance, updateId: 201, message: "starting now" }), + ); + const posts = calls.filter((c) => c.method === "post"); + expect(posts.length).toBe(3); + expect(posts[1]?.thread_ts).toBe("1700000000.0001"); + expect(posts[1]?.text).toContain("scheduled"); + expect(posts[2]?.thread_ts).toBe("1700000000.0001"); + expect(posts[2]?.text).toContain("starting now"); + expect(calls.filter((c) => c.method === "update").length).toBe(1); + }); + + test("a maintenance and a report sharing an id keep separate threads", async () => { + const { client, calls } = makeClient(); + const store = createMemoryAnchorStore(); + const channel = createSlackChannel({ + store, + createClient: () => client, + getBotToken: token, + softUnsubscribe: noopUnsub, + }); + + await channel.sendNotifications([makeSub()], makeUpdate({ updateId: 100 })); + await channel.sendNotifications( + [makeSub()], + makeUpdate({ status: "maintenance", updateId: 100 }), + ); + + // Both are roots: no thread reply, no root re-render. + expect(calls.filter((c) => c.method === "post").length).toBe(2); + expect(calls.every((c) => c.thread_ts === undefined)).toBe(true); expect(calls.some((c) => c.method === "update")).toBe(false); - expect(await store.getAnchor(10, 1)).toBeNull(); + expect((await store.getAnchor(REPORT, 1))?.ts).toBe("1700000000.0001"); + expect((await store.getAnchor(MAINTENANCE, 1))?.ts).toBe( + "1700000000.0002", + ); }); }); @@ -320,16 +364,16 @@ describe("validateSlackConfig", () => { describe("reserveDelivery (memory store)", () => { test("only the first reservation for a key wins", async () => { const store = createMemoryAnchorStore(); - expect(await store.reserveDelivery(1, 2, 3)).toBe(true); - expect(await store.reserveDelivery(1, 2, 3)).toBe(false); + expect(await store.reserveDelivery(REPORT, 2, 3)).toBe(true); + expect(await store.reserveDelivery(REPORT, 2, 3)).toBe(false); // A different updateId is an independent reservation. - expect(await store.reserveDelivery(1, 2, 4)).toBe(true); + expect(await store.reserveDelivery(REPORT, 2, 4)).toBe(true); }); test("releaseDelivery makes the key reservable again", async () => { const store = createMemoryAnchorStore(); - expect(await store.reserveDelivery(1, 2, 3)).toBe(true); - await store.releaseDelivery(1, 2, 3); - expect(await store.reserveDelivery(1, 2, 3)).toBe(true); + expect(await store.reserveDelivery(REPORT, 2, 3)).toBe(true); + await store.releaseDelivery(REPORT, 2, 3); + expect(await store.reserveDelivery(REPORT, 2, 3)).toBe(true); }); }); diff --git a/packages/subscriptions/src/channels/slack.ts b/packages/subscriptions/src/channels/slack.ts index 0c81d01cc..ef6208c0e 100644 --- a/packages/subscriptions/src/channels/slack.ts +++ b/packages/subscriptions/src/channels/slack.ts @@ -8,7 +8,11 @@ import { WebClient } from "@slack/web-api"; import type { PageUpdate, Subscription } from "../types"; import { buildReplyMessage, buildRootMessage } from "./slack-blocks"; -import { type SlackAnchorStore, createRedisAnchorStore } from "./slack-store"; +import { + type SlackAnchorStore, + type SlackThreadEvent, + createRedisAnchorStore, +} from "./slack-store"; interface SlackPostResult { ts?: string; @@ -124,40 +128,29 @@ export function createSlackChannel(deps: SlackChannelDeps) { } } - async function deliverMaintenance( + // Every event kind threads: the first delivered update becomes the root, + // later ones reply under it and re-render the root. + async function deliverThreaded( client: SlackClient, sub: Subscription, channelId: string, pageUpdate: PageUpdate, ): Promise { - const root = buildRootMessage(pageUpdate, sub); - const res = await runSlack(sub.id, () => - client.postMessage({ - channel: channelId, - attachments: root.attachments, - }), - ); - if (res === TEAM_TOKEN_INVALID) return TEAM_TOKEN_INVALID; - } - - async function deliverReport( - client: SlackClient, - sub: Subscription, - channelId: string, - pageUpdate: PageUpdate, - ): Promise { - const reportId = pageUpdate.id; + const event: SlackThreadEvent = { + kind: pageUpdate.status === "maintenance" ? "maintenance" : "report", + id: pageUpdate.id, + }; const updateId = pageUpdate.updateId; if (updateId == null) { - console.error(`slack: status report update ${reportId} missing updateId`); + console.error(`slack: ${event.kind} ${event.id} missing updateId`); return; } // Atomic dedupe: only the caller that wins this reservation posts. On a // failed post we release it below so the delivery stays retriable. - if (!(await deps.store.reserveDelivery(reportId, sub.id, updateId))) return; + if (!(await deps.store.reserveDelivery(event, sub.id, updateId))) return; - const anchor = await deps.store.getAnchor(reportId, sub.id); + const anchor = await deps.store.getAnchor(event, sub.id); const root = buildRootMessage(pageUpdate, sub); if (!anchor) { @@ -168,13 +161,13 @@ export function createSlackChannel(deps: SlackChannelDeps) { }), ); if (!res || res === TEAM_TOKEN_INVALID) { - await deps.store.releaseDelivery(reportId, sub.id, updateId); + await deps.store.releaseDelivery(event, sub.id, updateId); return res === TEAM_TOKEN_INVALID ? TEAM_TOKEN_INVALID : undefined; } if (res.ts) { // Stash this first update so the next one can backfill it into the // thread before the root is re-rendered and its content lost. - await deps.store.setAnchor(reportId, sub.id, { + await deps.store.setAnchor(event, sub.id, { ts: res.ts, channelId, pendingRootReply: buildReplyMessage(pageUpdate), @@ -196,14 +189,14 @@ export function createSlackChannel(deps: SlackChannelDeps) { }), ); if (!backfillRes || backfillRes === TEAM_TOKEN_INVALID) { - await deps.store.releaseDelivery(reportId, sub.id, updateId); + await deps.store.releaseDelivery(event, sub.id, updateId); return backfillRes === TEAM_TOKEN_INVALID ? TEAM_TOKEN_INVALID : undefined; } // Clear before posting the current reply: if that reply fails and the // delivery is retried, the first update must not be backfilled twice. - await deps.store.setAnchor(reportId, sub.id, { + await deps.store.setAnchor(event, sub.id, { ts: anchor.ts, channelId: anchor.channelId, }); @@ -219,7 +212,7 @@ export function createSlackChannel(deps: SlackChannelDeps) { }), ); if (!replyRes || replyRes === TEAM_TOKEN_INVALID) { - await deps.store.releaseDelivery(reportId, sub.id, updateId); + await deps.store.releaseDelivery(event, sub.id, updateId); return replyRes === TEAM_TOKEN_INVALID ? TEAM_TOKEN_INVALID : undefined; } @@ -266,10 +259,12 @@ export function createSlackChannel(deps: SlackChannelDeps) { // Sequential per team so a token failure aborts the batch before // hammering Slack with N calls that will all fail identically. for (const { sub, channelId } of members) { - const outcome = - pageUpdate.status === "maintenance" - ? await deliverMaintenance(client, sub, channelId, pageUpdate) - : await deliverReport(client, sub, channelId, pageUpdate); + const outcome = await deliverThreaded( + client, + sub, + channelId, + pageUpdate, + ); if (outcome === TEAM_TOKEN_INVALID) { console.error( `slack: team ${teamId} bot token invalid — aborting ${members.length} deliveries; subscribers left intact (reconnect the Slack app)`, diff --git a/packages/subscriptions/src/dispatcher.test.ts b/packages/subscriptions/src/dispatcher.test.ts index d17ea2624..9e9800a4a 100644 --- a/packages/subscriptions/src/dispatcher.test.ts +++ b/packages/subscriptions/src/dispatcher.test.ts @@ -29,6 +29,7 @@ import { import { assertSpyCalls, type Stub, stub } from "@std/testing/mock"; import { + dispatchMaintenance, dispatchMaintenanceUpdate, dispatchPageUpdate, dispatchStatusReportUpdate, @@ -243,6 +244,46 @@ describe("dispatchPageUpdate - edge cases", () => { }); }); +describe("dispatchMaintenance", () => { + test("announcement carries the first update's id and message", async () => { + const record = await db + .insert(maintenance) + .values({ + workspaceId: WORKSPACE_ID, + pageId: PAGE_ID, + title: "Database maintenance", + message: "stale column", + from: new Date("2026-08-10T10:00:00.000Z"), + to: new Date("2026-08-10T11:00:00.000Z"), + }) + .returning() + .get(); + + try { + const update = await db + .insert(maintenanceUpdate) + .values({ + maintenanceId: record.id, + message: "announcement", + date: new Date("2026-08-07T14:00:00.000Z"), + }) + .returning() + .get(); + + await dispatchMaintenance(record.id); + + const args = sendStatusReportUpdateMock.calls[0].args[0]; + expect(args.message).toBe("announcement"); + expect(args.date).toBe("2026-08-10T10:00:00.000Z"); + expect(args.idempotencyKey).toMatch( + new RegExp(`^maintenance-update:${update.id}:`), + ); + } finally { + await db.delete(maintenance).where(eq(maintenance.id, record.id)); + } + }); +}); + describe("dispatchMaintenanceUpdate", () => { test("dispatches the selected update with parent schedule and components", async () => { const startsAt = new Date("2026-08-10T10:00:00.000Z"); diff --git a/packages/subscriptions/src/dispatcher.ts b/packages/subscriptions/src/dispatcher.ts index caebca7de..d2c2c8918 100644 --- a/packages/subscriptions/src/dispatcher.ts +++ b/packages/subscriptions/src/dispatcher.ts @@ -127,6 +127,8 @@ export async function dispatchMaintenance(maintenanceId: number) { pageComponentIds: pageComponents.map((c) => c.id), pageComponents: pageComponents.map((c) => c.name), date: record.from.toISOString(), + // anchors the Slack thread and the email idempotency key on the first update + updateId: record.maintenanceUpdates[0]?.id, startsAt: record.from.toISOString(), endsAt: record.to.toISOString(), pageComponentsWithId: pageComponents.map((c) => ({