Something went wrong. Try again.
[READ-ONLY] Mirror of https://github.com/openstatusHQ/openstatus. ๐ซ Status page with uptime monitoring & API monitoring as code ๐ซ openstatus.dev
bun drizzle-orm monitoring monitoring-as-code nextjs observability on-call open-source shadcn-ui status-page statuspage synthetic-monitoring tinybird turso uptime uptime-checker uptime-monitor
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380import { getLogger } from "@logtape/logtape";import { isFeatureEnabled } from "@openstatus/services";import { getIncidentBySlackChannel } from "@openstatus/services/incident";import { missingSlackScopes, 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";import { type Block, buildAnswerMessage, buildConfirmationBlocks, buildLinkAccountBlocks, getConfirmationText, LINK_ACCOUNT_TEXT, type RefResolvers,} from "./blocks";import { channelContextTooling, contextChannelId, contextUserId, forgetContext, recallContext, rememberContext,} from "./channel-context";import type { SlackConfig, SlackEnv } from "./config";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, getStatusReportLink,} from "./page-urls";import { getRegistryTool, isSlackToolDraft, type SlackToolDraft,} from "./registry-runner";import { claimLinkCardWindow, claimPlanNoticeWindow, linkAccountUrl, planRequiredMessage, releaseLinkCardWindow, requireSlackMember, type SlackActor, slackAgentAllowed,} from "./require-slack-member";import { abortTurn, broadcastStop, endTurn, startTurn } from "./running-turns";import { buildThreadTitle, isThreadTitled, markThreadTitled, renameThread,} from "./thread-title";import { resolveWorkspace, type SlackWorkspace } from "./workspace-resolver";
const logger = getLogger("api-server");
const DEDUP_TTL_SECONDS = 10 * 60;
/** * 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<boolean> { 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; }}
const slackEventSchema = z.object({ type: z.string(), event: z .object({ type: z.string(), subtype: z.string().optional(), text: z.string().optional(), user: z.string().optional(), channel: z.string().optional(), channel_type: z.string().optional(), ts: z.string().optional(), 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({ bot: z.array(z.string()).optional(), oauth: z.array(z.string()).optional(), }) .optional(), assistant_thread: z .object({ channel_id: z.string(), thread_ts: z.string(), user_id: z.string().optional(), }) .optional(), // `app_context_changed` โ what the user has on screen. Absent entities // (`"context": {}`) mean they moved somewhere with nothing to track. context: z .object({ entities: z .array( z.object({ type: z.string().optional(), value: z.string().optional(), }), ) .optional(), }) .optional(), }) .optional(), event_id: z.string().optional(), team_id: z.string().optional(), challenge: z.string().optional(), // `app_context_changed` carries no `event.user`; the human is in here. authorizations: z .array( z.object({ user_id: z.string().optional(), is_bot: z.boolean().optional(), }), ) .optional(),});
type SlackEvent = z.infer<typeof slackEventSchema>;
const threadMessageSchema = z.object({ user: z.string().optional(), bot_id: z.string().optional(), text: z.string().optional(), ts: z.string().optional(),});
type ThreadMessage = z.infer<typeof threadMessageSchema>;
const slackPlatformErrorSchema = z.object({ code: z.literal("slack_webapi_platform_error"), data: z.object({ error: z.string(), }),});
function isSlackPlatformError(err: unknown, errorCode: string): boolean { const parsed = slackPlatformErrorSchema.safeParse(err); return parsed.success && parsed.data.data.error === errorCode;}
// Asking permission to run a write tool ("shall I go ahead and publish this?").const PERMISSION_QUESTION = /\b(shall i|should i|do you want me to|would you like me to|want me to|ready for me to|can i)\b[^?]*\b(go ahead|proceed|publish|create|post|send|schedule|submit|add|resolve|update)\b/i;
// A change written out as message text rather than passed to a tool.const PROSE_DRAFT_FIELDS = /\*{0,2}(title|status|message|impact|components?|from|to)\*{0,2}\s*:/gi;
/** * Heuristic for the "model drafted in prose instead of calling the tool" failure * โ logging only, so false positives are cheap. Both signals are required: a * plain answer can end in a question, and a list of reports can look tabular. */export function looksLikeUncardedDraft(text: string): boolean { if (!text) return false; if (!PERMISSION_QUESTION.test(text)) return false; return (text.match(PROSE_DRAFT_FIELDS) ?? []).length >= 2;}
// Bounds the Slack calls (and the agent's context) on very long threads.const MAX_THREAD_PAGES = 5;
// Replies come oldest first, so a single page would miss the latest messages// of a long thread โ the ones the agent is being asked about.async function fetchThread( slack: WebClient, channel: string, threadTs: string,): Promise<ThreadMessage[]> { const messages: ThreadMessage[] = []; let cursor: string | undefined; for (let page = 0; page < MAX_THREAD_PAGES; page++) { const replies = await slack.conversations.replies({ channel, ts: threadTs, limit: 100, cursor, }); messages.push(...((replies.messages ?? []) as ThreadMessage[])); cursor = replies.response_metadata?.next_cursor || undefined; if (!replies.has_more || !cursor) break; } return messages;}
/** * Whether an untagged channel-thread message is the user answering the agent * (e.g. "API" after "Which status page โ API or Marketing?"). True only when * the agent posted the message right before it and its author is whoever first * mentioned the agent in the thread. Anything else still needs a mention, so * the agent stays out of the humans' side of an incident thread. */export function isAnswerToAgent( thread: ThreadMessage[], message: { ts: string; user?: string }, botUserId: string,): boolean { if (!message.user || !botUserId) return false; const index = thread.findIndex((m) => m.ts === message.ts); const earlier = index === -1 ? thread : thread.slice(0, index);
const previous = earlier.at(-1); if (previous?.user !== botUserId) return false;
const starter = earlier.find( (m) => m.user !== botUserId && m.text?.includes(`<@${botUserId}>`), ); return starter?.user === message.user;}
export async function handleSlackEvent(c: Context<SlackEnv>) { const body = c.get("slackBody") as SlackEvent; const config = c.get("slackConfig");
if (body.type === "url_verification") { return c.json({ challenge: body.challenge }); }
if (body.type !== "event_callback") { return c.json({ ok: true }); }
if (body.event_id && (await isDuplicate(body.event_id))) { return c.json({ ok: true }); }
runInBackground("event", () => processEvent(body, config), { teamId: body.team_id, eventId: body.event_id, });
return c.json({ ok: true });}
/** Tells the agent which incident a channel belongs to, when it has one. */async function boundIncidentNote(args: { workspace: SlackWorkspace["workspace"]; actor: SlackActor; teamId: string; channelId: string | undefined;}): Promise<string | undefined> { if (!args.channelId) return undefined; const bound = await getIncidentBySlackChannel({ ctx: { workspace: args.workspace, actor: args.actor }, input: { teamId: args.teamId, channelId: args.channelId }, }).catch(() => undefined); if (!bound) return undefined; const report = bound.statusReport ? ` Its status report is "${bound.statusReport.title}" (id ${bound.statusReport.id}, ${bound.statusReport.status}).` : " It has no status report yet."; return `Incident channel: <#${args.channelId}> belongs to managed incident "${bound.title}" (id ${bound.id}, ${bound.severity}, ${bound.status}${bound.closedAt ? ", closed" : ""}).${report} Notes and status changes discussed here are about this incident: use id ${bound.id} without asking.`;}
/** * Tells an unlinked Slack user how to link their account. Passive surfaces * (mentions, DMs) send it at most once per window, so a chatty user isn't * flooded with cards. */async function sendLinkCard(args: { config: SlackConfig; workspaceId: number; teamId: string; slackUserId: string; post: (message: { text: string; blocks: Block[]; }) => Promise<{ ok?: boolean }>;}): Promise<void> { const { config, workspaceId, teamId, slackUserId, post } = args; if (!(await claimLinkCardWindow(teamId, slackUserId))) return; try { const url = await linkAccountUrl(config, { workspaceId, teamId, slackUserId, }); await post({ text: LINK_ACCOUNT_TEXT, blocks: buildLinkAccountBlocks(url), }); } catch (err) { // Otherwise a transient failure silences the card for the whole window. await releaseLinkCardWindow(teamId, slackUserId).catch(() => undefined); throw err; } logger.info("slack link card sent", { teamId, slackUserId });}
async function processEvent(body: SlackEvent, config: SlackConfig) { const event = body.event; if (!event) return;
if (event.type === "app_uninstalled" || event.type === "tokens_revoked") { const teamId = body.team_id; 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; }
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. if (event.type === "app_home_opened") { const tab = event.tab ?? "home"; if (tab !== "home" && tab !== "messages") return; const teamId = body.team_id; const userId = event.user; if (!teamId || !userId) return; const resolved = await resolveWorkspace(teamId); if (!resolved) return; const slack = new WebClient(resolved.botToken); const member = { workspace: resolved.workspace, teamId, slack };
if (tab === "messages") { const channel = event.channel; if (!channel) return; const actor = await requireSlackMember({ ...member, slackUserId: userId, }); if (!actor) { await sendLinkCard({ config, workspaceId: resolved.workspace.id, teamId, slackUserId: userId, post: (message) => slack.chat.postMessage({ channel, ...message }), }).catch((error) => logger.error("slack failed to send link card", { error, teamId }), ); return; } try { await greetOnce({ slack, teamId, userId, channel, }); } catch (err) { logger.error("slack failed to greet user", { error: err, teamId }); } return; }
try { const actor = await requireSlackMember({ ...member, slackUserId: userId, }); if (actor) { const needsReconnect = isFeatureEnabled(resolved.workspace, "incident-management") && missingSlackScopes(resolved.scopes).length > 0; await publishHomeView(slack, userId, { reconnectUrl: needsReconnect ? `${config.dashboardUrl}/settings/integrations` : undefined, }); } else { const url = await linkAccountUrl(config, { workspaceId: resolved.workspace.id, teamId, slackUserId: userId, }); await publishLinkAccountView(slack, userId, url); } } catch (err) { logger.error("slack failed to publish home view", { error: err, teamId }); } return; }
// Only fires while the app is still on `assistant_view`. Delete this branch // once the `agent_view` manifest is live โ the agent experience greets from // `app_home_opened` above instead. if (event.type === "assistant_thread_started") { const teamId = body.team_id; const thread = event.assistant_thread; if (!teamId || !thread?.user_id) return; const resolved = await resolveWorkspace(teamId); if (!resolved) return; const slack = new WebClient(resolved.botToken); const slackUserId = thread.user_id; try { const actor = await requireSlackMember({ workspace: resolved.workspace, teamId, slackUserId, slack, }); if (!actor) { await sendLinkCard({ config, workspaceId: resolved.workspace.id, teamId, slackUserId, post: (message) => slack.chat.postMessage({ channel: thread.channel_id, thread_ts: thread.thread_ts, ...message, }), }); return; } await greetOnce({ slack, teamId, userId: slackUserId, channel: thread.channel_id, threadTs: thread.thread_ts, }); } catch (err) { logger.error("slack failed to greet in new thread", { error: err, teamId, }); } return; }
// What the user is looking at, remembered for the next turn. Nothing is // read here โ this only records where they are. if (event.type === "app_context_changed") { const teamId = body.team_id; const userId = contextUserId(body.authorizations); if (!teamId || !userId) return;
const channelId = contextChannelId(event.context?.entities); if (channelId) { await rememberContext(teamId, userId, channelId); } else { await forgetContext(teamId, userId); } return; }
// A person renamed the thread, so the name is theirs now โ record it so no // later turn overwrites it. if (event.type === "agent_session_title_changed") { if (!body.team_id || !event.channel || !event.thread_ts) return; await markThreadTitled(body.team_id, event.channel, event.thread_ts); logger.info("slack thread renamed by user", { teamId: body.team_id, channel: event.channel, threadTs: event.thread_ts, }); return; }
// The user pressed stop. Slack has already halted any streamed message and // will not clear the session status itself, so both are on us. if (event.type === "agent_session_stopped") { const teamId = body.team_id; const channel = event.channel; const threadTs = event.thread_ts; 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, threadTs, wasRunning, });
// Cleared here as well as in the aborted turn's own cleanup: the turn may // be running on another instance, or already be gone. const resolved = await resolveWorkspace(teamId); if (!resolved) return; await setSessionStatus( new WebClient(resolved.botToken), channel, threadTs, "active", ).catch((err: unknown) => logger.warn("slack failed to clear status after stop", { error: err, channel, teamId, }), ); return; }
if (event.type !== "app_mention" && event.type !== "message") return; if (event.type === "message" && event.bot_id) return;
// The agent pane is the app's DM: every message there is addressed to us, // so no mention is required. const isAgentThread = event.channel_type === "im";
const ignoredSubtypes = [ "channel_join", "channel_leave", "channel_topic", "channel_purpose", "channel_name", ]; if (event.subtype && ignoredSubtypes.includes(event.subtype)) return; // In the agent pane, subtypes are the thread root (`assistant_app_thread`), // edits and deletions โ only plain user messages start a turn. if (isAgentThread && event.subtype) return;
const teamId = body.team_id; if (!teamId || !event.channel || !event.ts) return;
// 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. if (await isDuplicate(`msg:${event.channel}:${event.ts}`)) return;
const resolved = await resolveWorkspace(teamId); if (!resolved) { logger.warn("slack integration not found", { teamId }); return; }
const slack = new WebClient(resolved.botToken); const botUserId = resolved.botUserId; const threadTs = event.thread_ts ?? event.ts;
// Fetched early only when needed to decide whether to answer; reused below // so the agent sees the same thread. let prefetchedThread: ThreadMessage[] | undefined; if ( !isAgentThread && event.type === "message" && !event.text?.includes(`<@${botUserId}>`) ) { if (!event.thread_ts) return; try { prefetchedThread = await fetchThread( slack, event.channel, event.thread_ts, ); } catch (err) { logger.warn("slack failed to fetch thread for untagged reply", { error: err, channel: event.channel, teamId, }); return; } if ( !isAnswerToAgent( prefetchedThread, { ts: event.ts, user: event.user }, botUserId, ) ) { return; } }
logger.info("slack event received", { teamId, channel: event.channel, eventType: event.type, threadTs, user: event.user, agentThread: isAgentThread, });
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, slackUserId, slack, }); if (!actor) { await sendLinkCard({ config, workspaceId: resolved.workspace.id, teamId, slackUserId, post: (message) => 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 link card", { error, teamId }), ); return; }
// Registered before the session is marked `processing`: a stop landing in // that window has to find a controller, or the turn runs on unstoppable. const turn = startTurn(event.channel, threadTs); let reply: Reply | undefined;
try { reply = await createReply({ slack, channel: event.channel, threadTs, teamId, userId: event.user, isAgentThread, });
if (turn.signal.aborted) { logger.info("slack turn stopped before it started", { teamId, channel: event.channel, threadTs, }); await reply.stopped(); return; }
let thread: ThreadMessage[] = []; if (prefetchedThread) { thread = prefetchedThread; } else if (event.thread_ts) { thread = await fetchThread(slack, event.channel, event.thread_ts); } else { thread = [{ user: event.user, text: event.text, ts: event.ts }]; }
// Only in the pane: in a channel the agent is already reading the thread // it was mentioned in, and "what you're looking at" is that same channel. const contextChannel = isAgentThread && event.user ? await recallContext(teamId, event.user) : undefined; const context = contextChannel ? channelContextTooling({ slack, channelId: contextChannel }) : undefined; const incidentNote = await boundIncidentNote({ workspace: resolved.workspace, actor, teamId, channelId: isAgentThread ? contextChannel : event.channel, });
logger.info("slack agent invoked", { teamId, channel: event.channel, threadTs, messageCount: thread.length, contextChannel, });
const result = await runAgent( resolved.workspace, thread, botUserId, event.text, actor, { events: reply.progress, signal: turn.signal, tools: context?.tools, contextNote: [context?.contextNote, incidentNote].filter(Boolean).join("\n\n") || undefined, }, );
if (result.aborted) { logger.info("slack turn abandoned after stop", { teamId, channel: event.channel, threadTs, }); await reply.stopped(); return; }
logger.info("slack agent completed", { teamId, channel: event.channel, threadTs, toolCalls: result.toolResults.map((tr) => tr.toolName), finishReason: result.finishReason, stepCount: result.stepCount, hitStepLimit: result.hitStepLimit, });
// One card per draft. A turn can draft several changes (rename a report // *and* post an update to it); dropping all but the first would leave the // user with nothing to click for the rest. When the model drafts the same // thing twice, its last draft is the one it meant. const drafts = new Map<string, SlackToolDraft>(); for (const tr of result.toolResults) { if (!isSlackToolDraft(tr.result)) continue; const key = draftKey({ toolName: tr.result.toolName, input: tr.result.input, }); drafts.delete(key); drafts.set(key, tr.result); } const firstDraft = drafts.values().next().value;
if (firstDraft) { for (const draft of drafts.values()) { logger.info("slack confirmation requested", { teamId, channel: event.channel, threadTs, toolName: draft.toolName, }); await handleConfirmation( slack, reply, event.channel, threadTs, event.user ?? "", resolved.workspace.id, teamId, draft, ); } } else { // No draft means no card. Distinguish a legitimate text answer from the // model drafting a change in prose and asking for permission instead of // calling the tool โ the latter strands the user with nothing to click. if (looksLikeUncardedDraft(result.text)) { logger.warn("slack draft proposed without a card", { teamId, channel: event.channel, threadTs, finishReason: result.finishReason, stepCount: result.stepCount, hitStepLimit: result.hitStepLimit, readToolCalls: result.toolResults.map((tr) => tr.toolName), }); } if (result.text) { await reply.answer(result.text); } else { await reply.send({ text: "Done!" }); } logger.info("slack response sent", { teamId, channel: event.channel, threadTs, }); }
await titleThread({ slack, channel: event.channel, threadTs, teamId, workspaceId: resolved.workspace.id, isAgentThread, draft: firstDraft, userText: event.text, }); } catch (err) { logger.error("slack agent error", { error: err, channel: event.channel, teamId, threadTs, }); await reply ?.send({ text: ":x: Something went wrong. Please try again." }) .catch((sendErr: unknown) => { logger.error("slack failed to send error message", { error: sendErr, channel: event.channel, threadTs, }); }); } finally { endTurn(event.channel, threadTs, turn); await reply?.finish?.(); }}
function statusReportIdOf(input: unknown): number | undefined { if (typeof input !== "object" || input === null) return undefined; const id = (input as { statusReportId?: unknown }).statusReportId; return typeof id === "number" ? id : undefined;}
/** * Names the thread after its subject, so the agent pane's timeline reads as a * list of incidents rather than a list of operations. * * Confined to the agent pane: that timeline is the whole payoff, and Slack * documents `agents.sessions.rename` as also renaming the channel for session * channels โ not something to risk on a shared incident channel. * * Cosmetic, and last: a failure here must never cost the user their answer. */async function titleThread(args: { slack: WebClient; channel: string; threadTs: string; teamId: string; workspaceId: number; isAgentThread: boolean; draft?: { toolName: string; input: unknown }; userText?: string;}): Promise<void> { const { slack, channel, threadTs, teamId, workspaceId, isAgentThread, draft, userText, } = args; if (!isAgentThread) return;
try { if (await isThreadTitled(teamId, channel, threadTs)) return;
// Only looked up once, and only for a draft that acts on an existing // report โ `add_update` and `resolve` carry an id but no title. let reportTitle: string | undefined; const reportId = draft && statusReportIdOf(draft.input); if (reportId !== undefined) { const link = await getStatusReportLink(workspaceId, reportId); reportTitle = link?.title; }
const title = buildThreadTitle({ draft, reportTitle, userText }); if (!title) return;
await renameThread({ slack, channel, threadTs, title, teamId }); } catch (err) { logger.warn("slack failed to title the thread", { error: err, channel, teamId, }); }}
/** * 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; or a native "working" status on the thread. */interface Reply { /** * Delivers the agent's free-text answer and returns the ts of the message * holding it. When the answer was streamed it is already on screen, and this * only finalizes the message. */ answer(text: string): Promise<string>; /** Writes a message of our own โ a confirmation card, or an error. */ send(message: { text: string; blocks?: Block[] }): Promise<string>; /** * Confirms the turn ended because the user stopped it, leaving whatever was * already written in place. */ stopped(): Promise<void>; /** Progress to report while the turn runs, when the surface can show it. */ progress?: AgentEvents; /** Runs once the turn is over, whether it succeeded or not. */ finish?: () => Promise<void>;}
const STOPPED_NOTICE = "_Stopped._";
/** * Registry tools are named `verb_noun` (`list_status_pages`, * `get_monitor_status`). A task list reads better as an activity, so the verb * becomes a gerund and the rest is left as words. */const TOOL_ACTIVITY: Record<string, string> = { list: "Reading", get: "Reading", search: "Searching", create: "Drafting", add: "Drafting", update: "Drafting", resolve: "Drafting",};
export function toolTaskTitle(toolName: string): string { const [verb, ...rest] = toolName.split("_"); const subject = rest.join(" "); const activity = TOOL_ACTIVITY[verb]; if (!activity || !subject) return toolName.replace(/_/g, " "); return `${activity} ${subject}`;}
function taskChunk( id: string, toolName: string, status: "in_progress" | "complete",) { return { type: "task_update" as const, id, title: toolTaskTitle(toolName), status, };}
/** * Opens a stream for the turn, or returns undefined when this surface can't * carry one. Whether the *workspace* allows streaming only shows up on the * first append, mid-turn โ `streamingReply` handles that failure. */function createStreamer(args: { slack: WebClient; channel: string; threadTs: string; teamId: string; userId: string | undefined; isAgentThread: boolean;}) { const { slack, channel, threadTs, teamId, userId, isAgentThread } = args; // Older @slack/web-api has no streaming support. if (typeof slack.chatStream !== "function") return undefined; // Outside a DM, Slack needs to know who the streamed message is for. if (!isAgentThread && !userId) return undefined; try { return slack.chatStream({ channel, thread_ts: threadTs, task_display_mode: "timeline", recipient_user_id: userId, recipient_team_id: teamId, }); } catch (err) { logger.warn("slack could not open a stream", { error: err, channel, teamId, }); return undefined; }}
/** * Streams the answer as the model writes it. Every stream call is best-effort: * a workspace without streaming enabled only fails on the first append, by * which point the turn is already running, so a failure latches into `broken` * and the answer is posted (or the half-written message rewritten) instead. */function streamingReply(args: { slack: WebClient; streamer: NonNullable<ReturnType<typeof createStreamer>>; channel: string; threadTs: string; teamId: string; finishSession?: () => Promise<void>;}): Reply { const { slack, streamer, channel, threadTs, teamId, finishSession } = args; const post = postInThread(slack, channel, threadTs);
let appended = false; let streamedText = false; let broken = false; let streamClosed = false;
const attempt = async (fn: () => Promise<void>) => { if (broken) return; try { await fn(); appended = true; } catch (err) { broken = true; logger.warn("slack stream failed, falling back to a posted message", { error: err, channel, teamId, }); } };
/** * Finalizes the streamed message, if one was ever opened โ including after a * later append failed, or Slack leaves that message streaming for good. */ const closeStream = async (): Promise<string | undefined> => { if (streamClosed || !appended) return streamer.ts; streamClosed = true; try { await streamer.stop(); } catch (err) { broken = true; logger.warn("slack failed to stop the stream", { error: err, channel, teamId, }); } return streamer.ts; };
return { progress: { onTextDelta: (delta) => attempt(async () => { await streamer.append({ markdown_text: delta }); streamedText = true; }), onToolCall: ({ id, toolName }) => attempt(async () => { await streamer.append({ chunks: [taskChunk(id, toolName, "in_progress")], }); }), onToolResult: ({ id, toolName }) => attempt(async () => { await streamer.append({ chunks: [taskChunk(id, toolName, "complete")], }); }), }, async answer(text) { const ts = await closeStream(); if (streamedText && !broken) return ts ?? ""; // The stream never carried the answer โ nothing was streamed, or it // broke partway. Put the whole answer on screen, rewriting the // half-written message when there is one. const message = buildAnswerMessage(text); if (ts) { await slack.chat.update({ channel, ts, ...message }); return ts; } return post(message); }, async send(message) { // The card is a message of its own: updating the streamed one in place // would wipe the answer Slack has already rendered. await closeStream(); return post(message); }, async stopped() { if (streamClosed || !appended) { await post({ text: STOPPED_NOTICE }); return; } streamClosed = true; try { // Closes the partial answer with the notice attached. Slack halts the // stream on its side when the user presses stop, so this often fails โ // the partial message is already final, which is the point. await streamer.stop({ markdown_text: `\n\n${STOPPED_NOTICE}` }); } catch (err) { logger.info("slack stream already closed by the stop request", { error: err, channel, teamId, }); } }, async finish() { await closeStream(); await finishSession?.(); }, };}
/** * Picks the richest delivery this surface supports. The session status is * independent of streaming โ a streamed turn still marks the thread as * working, so Slack shows the loading state and the stop button. */async function createReply(args: { slack: WebClient; channel: string; threadTs: string; teamId: string; userId: string | undefined; isAgentThread: boolean;}): Promise<Reply> { const { slack, channel, threadTs, teamId, userId } = args;
const session = await acknowledgeWithSession( slack, channel, threadTs, teamId, userId, );
// 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) => logger.warn("slack failed to set assistant status", { error: err, channel, teamId, }), ); }
const streamer = createStreamer(args); if (streamer) { return streamingReply({ slack, streamer, channel, threadTs, teamId, finishSession: session?.finish, }); }
return session ?? acknowledgeWithStatus(slack, channel, threadTs);}
function postInThread( slack: WebClient, channel: string, threadTs: string,): Reply["send"] { return async ({ text, blocks }) => { try { const res = await slack.chat.postMessage({ channel, thread_ts: threadTs, text, blocks, }); if (!res.ts) throw new Error("chat.postMessage returned no ts"); return res.ts; } catch (err) { if (!isSlackPlatformError(err, "cannot_reply_to_message")) throw err; // The parent message can't host a thread โ answer at the top level // rather than losing the reply entirely. logger.warn("slack cannot reply to message, falling back to top-level", { channel, threadTs, }); const res = await slack.chat.postMessage({ channel, text, blocks }); if (!res.ts) throw new Error("chat.postMessage returned no ts"); return res.ts; } };}
/** * Preferred acknowledgement: mark the thread's agent session as `processing` * so Slack shows the agent working โ no placeholder message โ and hand it * back as `active` once we've answered. Returns undefined when the workspace * doesn't support agent sessions, so the caller falls back to the older * indicators. */async function acknowledgeWithSession( slack: WebClient, channel: string, threadTs: string, teamId: string, userId: string | undefined,): Promise<Reply | undefined> { try { await setSessionStatus(slack, channel, threadTs, "processing", userId); } catch (err) { logger.info("slack agent session unavailable, falling back", { error: err, channel, teamId, }); return; } const send = postInThread(slack, channel, threadTs); return { send, answer: (text) => send(buildAnswerMessage(text)), stopped: async () => { await send({ text: STOPPED_NOTICE }); }, async finish() { await setSessionStatus(slack, channel, threadTs, "active").catch( (err: unknown) => logger.warn("slack failed to reset agent session status", { error: err, channel, teamId, }), ); }, };}
// 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 acknowledgeWithStatus( slack: WebClient, channel: string, threadTs: string,): Reply { const send = postInThread(slack, channel, threadTs); return { send, answer: (text) => send(buildAnswerMessage(text)), stopped: async () => { await send({ text: STOPPED_NOTICE }); }, };}
async function handleConfirmation( slack: WebClient, reply: Reply, channel: string, threadTs: string, userId: string, workspaceId: number, teamId: string, draft: SlackToolDraft,) { const tool = getRegistryTool(draft.toolName); if (!tool) { logger.error("slack: registry tool not found", { toolName: draft.toolName, }); await reply.send({ text: ":x: Something went wrong. Please try again." }); return; }
const payload: PendingPayload = { toolName: draft.toolName, input: draft.input, }; const text = getConfirmationText({ tool, input: draft.displayInput });
// findByThread + replace isn't atomic on its own โ two concurrent // events on the same thread could both see `existing` and race on // 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);
const blocks = await buildConfirmationBlocks({ actionId: existing.id, tool, input: draft.displayInput, resolvers: makeRefResolvers(workspaceId), }); await reply.send({ text, blocks }); await slack.chat.update({ channel, ts: existing.messageTs, text, blocks, }); } else { // The card's buttons carry the action id, and the stored action carries // the card's ts โ so write the text first to learn the ts, then attach // the buttons. const messageTs = await reply.send({ text }); const actionId = await store({ workspaceId, teamId, channelId: channel, threadTs, messageTs, userId, payload, });
const blocks = await buildConfirmationBlocks({ actionId, tool, input: draft.displayInput, resolvers: makeRefResolvers(workspaceId), }); await slack.chat.update({ channel, ts: messageTs, text, blocks }); }}