diff --git a/apps/server/src/routes/slack/handler.test.ts b/apps/server/src/routes/slack/handler.test.ts index af44d067..ee0f8988 100644 --- a/apps/server/src/routes/slack/handler.test.ts +++ b/apps/server/src/routes/slack/handler.test.ts @@ -511,6 +511,68 @@ describe("handleSlackEvent", () => { expect(posts.length).toBe(0); }); + test("clears the status only after the last concurrent run in a thread", async () => { + // Slack keeps one status per (channel, thread); the first run to finish + // must not pull the indicator out from under the second. + let release: (() => void) | undefined; + const gate = new Promise((r) => { + release = r; + }); + let call = 0; + slackTestState.runAgentOverride = () => { + call += 1; + const answer = { text: "Here is my response", toolResults: [] }; + // The first run holds the lease until the second has come and gone. + return call === 1 ? gate.then(() => answer) : Promise.resolve(answer); + }; + + const threadTs = "700.10"; + const post = (ts: string) => + signAndPost(app, { + type: "event_callback", + team_id: "T_KNOWN", + event_id: `evt_concurrent_${ts}`, + event: { + type: "app_mention", + text: "<@UBOT> hello", + user: "U1", + channel: "C1", + ts, + thread_ts: threadTs, + }, + }); + + await post("700.11"); + await new Promise((r) => setTimeout(r, 50)); + await post("700.12"); + await new Promise((r) => setTimeout(r, 100)); + + // Run two has replied and released its hold; run one is still going, so + // the indicator must survive. + expect( + slackTestState.calls.filter((m) => m.method === "postMessage").length, + ).toBe(1); + expect( + slackTestState.calls.filter( + (m) => m.method === "setStatus" && m.args.status === "", + ).length, + ).toBe(0); + + release?.(); + await new Promise((r) => setTimeout(r, 100)); + + // Set once for the thread, cleared once, after both runs finished. + const statuses = slackTestState.calls.filter( + (m) => m.method === "setStatus", + ); + expect(statuses.length).toBe(2); + expect(statuses[0].args.status).toBe("is thinking..."); + expect(statuses[1].args.status).toBe(""); + expect( + slackTestState.calls.filter((m) => m.method === "postMessage").length, + ).toBe(2); + }); + test("sets and clears the thinking status around the agent run", async () => { const res = await signAndPost(app, { type: "event_callback", @@ -558,7 +620,8 @@ describe("handleSlackEvent", () => { }); test("still replies when setStatus fails", async () => { - slackTestState.setStatusOverride = () => { + slackTestState.setStatusOverride = (args: Record) => { + slackTestState.calls.push({ method: "setStatus", args }); const err = new Error("An API error occurred: channel_not_found"); Object.assign(err, { code: "slack_webapi_platform_error", @@ -567,26 +630,37 @@ describe("handleSlackEvent", () => { return Promise.reject(err); }; + const ts = `${Date.now()}.22`; const res = await signAndPost(app, { type: "event_callback", team_id: "T_KNOWN", - event_id: `evt_statusfail_${Date.now()}`, + event_id: `evt_statusfail_${ts}`, event: { type: "app_mention", text: "<@UBOT> hello", user: "U1", channel: "C1", - ts: `${Date.now()}.22`, + ts, }, }); expect(res.status).toBe(200); await new Promise((r) => setTimeout(r, 100)); + // The agent's answer, not the error branch's message. const posts = slackTestState.calls.filter( (m) => m.method === "postMessage", ); expect(posts.length).toBe(1); + expect(posts[0].args.text).toBe("Here is my response"); + expect(posts[0].args.thread_ts).toBe(ts); + + // A failing set must not skip the clear. + const statuses = slackTestState.calls.filter( + (m) => m.method === "setStatus", + ); + expect(statuses.length).toBe(2); + expect(statuses[1].args.status).toBe(""); }); test("shows error message when runAgent throws", async () => { diff --git a/apps/server/src/routes/slack/handler.ts b/apps/server/src/routes/slack/handler.ts index ffeb3927..dd0e0f50 100644 --- a/apps/server/src/routes/slack/handler.ts +++ b/apps/server/src/routes/slack/handler.ts @@ -311,9 +311,19 @@ async function processEvent(body: SlackEvent) { confirmationResult, ); } else { - await postThreadReply(slack, event.channel, threadTs, { + const posted = await postThreadReply(slack, event.channel, threadTs, { text: result.text || "Done!", }); + if (!posted) { + // postThreadReply already exhausted the top-level fallback, so the + // answer is lost; don't claim it was delivered. + logger.error("slack response not delivered", { + teamId, + channel: event.channel, + threadTs, + }); + return; + } logger.info("slack response sent", { teamId, channel: event.channel, @@ -393,7 +403,16 @@ async function handleConfirmation( text, blocks, }); - if (!messageTs) return; + if (!messageTs) { + // Without a card there is nothing to approve, so leave the pending record + // alone and let its TTL expire rather than pointing it at a dead message. + logger.error("slack confirmation card not delivered", { + channel, + threadTs, + toolName: draft.toolName, + }); + return; + } if (existing) { await replace(existing.id, payload, messageTs); diff --git a/apps/server/src/routes/slack/status.ts b/apps/server/src/routes/slack/status.ts index e129004f..4b6bc163 100644 --- a/apps/server/src/routes/slack/status.ts +++ b/apps/server/src/routes/slack/status.ts @@ -15,8 +15,19 @@ const LOADING_MESSAGES = [ "Drafting an update...", ]; +interface Lease { + refs: number; + timer?: ReturnType; + /** Serializes setStatus calls so the clear can't overtake a live refresh. */ + queue: Promise; +} + +// Slack stores one status per (channel, thread), so concurrent runs on the +// same thread share a lease and only the last one to finish clears it. +const leases = new Map(); + export interface ThinkingStatus { - /** Idempotent — clears the refresh timer and removes the indicator. */ + /** Idempotent — releases this run's hold on the thread's indicator. */ stop(): Promise; } @@ -30,38 +41,59 @@ export async function startThinkingStatus( channelId: string, threadTs: string, ): Promise { - const setStatus = (status: string, loadingMessages?: string[]) => - slack.assistant.threads - .setStatus({ - channel_id: channelId, - thread_ts: threadTs, - status, - ...(loadingMessages ? { loading_messages: loadingMessages } : {}), - }) - .then(() => undefined) - .catch((error: unknown) => { - logger.warn("slack failed to set thinking status", { - error, - channel: channelId, - threadTs, - }); - }); + const key = `${channelId}:${threadTs}`; - await setStatus(STATUS, LOADING_MESSAGES); + const send = (lease: Lease, status: string, loadingMessages?: string[]) => { + lease.queue = lease.queue.then(() => + slack.assistant.threads + .setStatus({ + channel_id: channelId, + thread_ts: threadTs, + status, + ...(loadingMessages ? { loading_messages: loadingMessages } : {}), + }) + .then(() => undefined) + .catch((error: unknown) => { + logger.warn("slack failed to set thinking status", { + error, + channel: channelId, + threadTs, + }); + }), + ); + return lease.queue; + }; + let lease = leases.get(key); + if (lease) { + lease.refs += 1; + } else { + lease = { refs: 1, queue: Promise.resolve() }; + leases.set(key, lease); + const held = lease; + held.timer = setInterval(() => { + void send(held, STATUS, LOADING_MESSAGES); + }, REFRESH_MS); + await send(held, STATUS, LOADING_MESSAGES); + } + + const held = lease; let stopped = false; - const timer = setInterval(() => { - void setStatus(STATUS, LOADING_MESSAGES); - }, REFRESH_MS); return { async stop() { if (stopped) return; stopped = true; - clearInterval(timer); - // Posting in-thread clears the indicator on its own, but the - // top-level fallback in the handler posts outside the thread. - await setStatus(""); + + held.refs -= 1; + if (held.refs > 0) return; + + clearInterval(held.timer); + leases.delete(key); + // Posting in-thread clears the indicator on its own, but the top-level + // fallback in the handler posts outside the thread. Queued last, so any + // refresh already in flight resolves before the clear is sent. + await send(held, ""); }, }; }