Skip to content

Commit c73e80b

Browse files
committed
fix(slack): make inbound retries explicit
1 parent bfc77b0 commit c73e80b

5 files changed

Lines changed: 116 additions & 40 deletions

File tree

extensions/slack/src/monitor/context.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ export type SlackMonitorContext = {
6868

6969
logger: ReturnType<typeof getChildLogger>;
7070
markMessageSeen: (channelId: string | undefined, ts?: string) => boolean;
71+
releaseSeenMessage: (channelId: string | undefined, ts?: string) => void;
7172
shouldDropMismatchedSlackEvent: (body: unknown) => boolean;
7273
resolveSlackSystemEventSessionKey: (params: {
7374
channelId?: string | null;
@@ -160,6 +161,13 @@ export function createSlackMonitorContext(params: {
160161
return seenMessages.check(`${channelId}:${ts}`);
161162
};
162163

164+
const releaseSeenMessage = (channelId: string | undefined, ts?: string) => {
165+
if (!channelId || !ts) {
166+
return;
167+
}
168+
seenMessages.delete(`${channelId}:${ts}`);
169+
};
170+
163171
const resolveSlackSystemEventSessionKey = (p: {
164172
channelId?: string | null;
165173
channelType?: string | null;
@@ -433,6 +441,7 @@ export function createSlackMonitorContext(params: {
433441
removeAckAfterReply: params.removeAckAfterReply,
434442
logger,
435443
markMessageSeen,
444+
releaseSeenMessage,
436445
shouldDropMismatchedSlackEvent,
437446
resolveSlackSystemEventSessionKey,
438447
isChannelAllowed,

extensions/slack/src/monitor/message-handler.app-mention-race.test.ts

Lines changed: 53 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -51,30 +51,41 @@ vi.mock("./message-handler/dispatch.js", () => ({
5151
}));
5252

5353
let createSlackMessageHandler: typeof import("./message-handler.js").createSlackMessageHandler;
54+
let SlackRetryableInboundError: typeof import("./message-handler.js").SlackRetryableInboundError;
5455

5556
function createMarkMessageSeen() {
5657
const seen = new Set<string>();
57-
return (channel: string | undefined, ts: string | undefined) => {
58-
if (!channel || !ts) {
58+
return {
59+
markMessageSeen(channel: string | undefined, ts: string | undefined) {
60+
if (!channel || !ts) {
61+
return false;
62+
}
63+
const key = `${channel}:${ts}`;
64+
if (seen.has(key)) {
65+
return true;
66+
}
67+
seen.add(key);
5968
return false;
60-
}
61-
const key = `${channel}:${ts}`;
62-
if (seen.has(key)) {
63-
return true;
64-
}
65-
seen.add(key);
66-
return false;
69+
},
70+
releaseSeenMessage(channel: string | undefined, ts: string | undefined) {
71+
if (!channel || !ts) {
72+
return;
73+
}
74+
seen.delete(`${channel}:${ts}`);
75+
},
6776
};
6877
}
6978

7079
function createTestHandler() {
80+
const seenMessages = createMarkMessageSeen();
7181
return createSlackMessageHandler({
7282
ctx: {
7383
cfg: {},
7484
accountId: "default",
7585
app: { client: {} },
7686
runtime: {},
77-
markMessageSeen: createMarkMessageSeen(),
87+
markMessageSeen: seenMessages.markMessageSeen,
88+
releaseSeenMessage: seenMessages.releaseSeenMessage,
7889
} as Parameters<typeof createSlackMessageHandler>[0]["ctx"],
7990
account: { accountId: "default" } as Parameters<typeof createSlackMessageHandler>[0]["account"],
8091
});
@@ -118,7 +129,8 @@ async function createInFlightMessageScenario(ts: string) {
118129

119130
describe("createSlackMessageHandler app_mention race handling", () => {
120131
beforeAll(async () => {
121-
({ createSlackMessageHandler } = await import("./message-handler.js"));
132+
({ createSlackMessageHandler, SlackRetryableInboundError } =
133+
await import("./message-handler.js"));
122134
});
123135

124136
beforeEach(() => {
@@ -183,4 +195,34 @@ describe("createSlackMessageHandler app_mention race handling", () => {
183195
expect(prepareSlackMessageMock).toHaveBeenCalledTimes(1);
184196
expect(dispatchPreparedSlackMessageMock).toHaveBeenCalledTimes(1);
185197
});
198+
199+
it("retries message replay after an explicit retryable dispatch failure", async () => {
200+
prepareSlackMessageMock.mockResolvedValue({ ctxPayload: {} });
201+
dispatchPreparedSlackMessageMock
202+
.mockRejectedValueOnce(new SlackRetryableInboundError("retry me"))
203+
.mockResolvedValueOnce(undefined);
204+
205+
const handler = createTestHandler();
206+
207+
await expect(sendMessageEvent(handler, "1700000000.000250")).rejects.toThrow("retry me");
208+
await expect(sendMessageEvent(handler, "1700000000.000250")).resolves.toBeUndefined();
209+
210+
expect(prepareSlackMessageMock).toHaveBeenCalledTimes(2);
211+
expect(dispatchPreparedSlackMessageMock).toHaveBeenCalledTimes(2);
212+
});
213+
214+
it("keeps message replay deduped after a non-retryable dispatch failure", async () => {
215+
prepareSlackMessageMock.mockResolvedValue({ ctxPayload: {} });
216+
dispatchPreparedSlackMessageMock.mockRejectedValueOnce(new Error("post-send failure"));
217+
218+
const handler = createTestHandler();
219+
220+
await expect(sendMessageEvent(handler, "1700000000.000300")).rejects.toThrow(
221+
"post-send failure",
222+
);
223+
await sendMessageEvent(handler, "1700000000.000300");
224+
225+
expect(prepareSlackMessageMock).toHaveBeenCalledTimes(1);
226+
expect(dispatchPreparedSlackMessageMock).toHaveBeenCalledTimes(1);
227+
});
186228
});

extensions/slack/src/monitor/message-handler.test.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ vi.mock("./thread-resolution.js", () => ({
3232

3333
function createContext(overrides?: {
3434
markMessageSeen?: (channel: string | undefined, ts: string | undefined) => boolean;
35+
releaseSeenMessage?: (channel: string | undefined, ts: string | undefined) => void;
3536
}) {
3637
return {
3738
cfg: {},
@@ -42,11 +43,14 @@ function createContext(overrides?: {
4243
runtime: {},
4344
markMessageSeen: (channel: string | undefined, ts: string | undefined) =>
4445
overrides?.markMessageSeen?.(channel, ts) ?? false,
46+
releaseSeenMessage: (channel: string | undefined, ts: string | undefined) =>
47+
overrides?.releaseSeenMessage?.(channel, ts),
4548
} as Parameters<typeof createSlackMessageHandler>[0]["ctx"];
4649
}
4750

4851
function createHandlerWithTracker(overrides?: {
4952
markMessageSeen?: (channel: string | undefined, ts: string | undefined) => boolean;
53+
releaseSeenMessage?: (channel: string | undefined, ts: string | undefined) => void;
5054
}) {
5155
const trackEvent = vi.fn();
5256
const handler = createSlackMessageHandler({

extensions/slack/src/monitor/message-handler.ts

Lines changed: 49 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ export type SlackMessageHandler = (
1717

1818
const APP_MENTION_RETRY_TTL_MS = 60_000;
1919

20+
export class SlackRetryableInboundError extends Error {
21+
constructor(message: string, options?: ErrorOptions) {
22+
super(message, options);
23+
this.name = "SlackRetryableInboundError";
24+
}
25+
}
26+
2027
function resolveSlackSenderId(message: SlackMessageEvent): string | null {
2128
return message.user ?? message.bot_id ?? null;
2229
}
@@ -133,40 +140,53 @@ export function createSlackMessageHandler(params: {
133140
...last.message,
134141
text: combinedText,
135142
};
136-
const prepared = await prepareSlackMessage({
137-
ctx,
138-
account,
139-
message: syntheticMessage,
140-
opts: {
141-
...last.opts,
142-
wasMentioned: combinedMentioned || last.opts.wasMentioned,
143-
},
144-
});
145143
const seenMessageKey = buildSeenMessageKey(last.message.channel, last.message.ts);
146-
if (!prepared) {
147-
return;
148-
}
149-
if (seenMessageKey) {
150-
pruneAppMentionRetryKeys(Date.now());
151-
if (last.opts.source === "app_mention") {
152-
// If app_mention wins the race and dispatches first, drop the later message dispatch.
153-
appMentionDispatchedKeys.set(seenMessageKey, Date.now() + APP_MENTION_RETRY_TTL_MS);
154-
} else if (last.opts.source === "message" && appMentionDispatchedKeys.has(seenMessageKey)) {
155-
appMentionDispatchedKeys.delete(seenMessageKey);
156-
appMentionRetryKeys.delete(seenMessageKey);
144+
try {
145+
const prepared = await prepareSlackMessage({
146+
ctx,
147+
account,
148+
message: syntheticMessage,
149+
opts: {
150+
...last.opts,
151+
wasMentioned: combinedMentioned || last.opts.wasMentioned,
152+
},
153+
});
154+
if (!prepared) {
157155
return;
158156
}
159-
appMentionRetryKeys.delete(seenMessageKey);
160-
}
161-
if (entries.length > 1) {
162-
const ids = entries.map((entry) => entry.message.ts).filter(Boolean) as string[];
163-
if (ids.length > 0) {
164-
prepared.ctxPayload.MessageSids = ids;
165-
prepared.ctxPayload.MessageSidFirst = ids[0];
166-
prepared.ctxPayload.MessageSidLast = ids[ids.length - 1];
157+
if (seenMessageKey) {
158+
pruneAppMentionRetryKeys(Date.now());
159+
if (last.opts.source === "app_mention") {
160+
// If app_mention wins the race and dispatches first, drop the later message dispatch.
161+
appMentionDispatchedKeys.set(seenMessageKey, Date.now() + APP_MENTION_RETRY_TTL_MS);
162+
} else if (
163+
last.opts.source === "message" &&
164+
appMentionDispatchedKeys.has(seenMessageKey)
165+
) {
166+
appMentionDispatchedKeys.delete(seenMessageKey);
167+
appMentionRetryKeys.delete(seenMessageKey);
168+
return;
169+
}
170+
appMentionRetryKeys.delete(seenMessageKey);
171+
}
172+
if (entries.length > 1) {
173+
const ids = entries.map((entry) => entry.message.ts).filter(Boolean) as string[];
174+
if (ids.length > 0) {
175+
prepared.ctxPayload.MessageSids = ids;
176+
prepared.ctxPayload.MessageSidFirst = ids[0];
177+
prepared.ctxPayload.MessageSidLast = ids[ids.length - 1];
178+
}
179+
}
180+
await dispatchPreparedSlackMessage(prepared);
181+
} catch (error) {
182+
if (error instanceof SlackRetryableInboundError) {
183+
if (seenMessageKey) {
184+
appMentionDispatchedKeys.delete(seenMessageKey);
185+
}
186+
ctx.releaseSeenMessage(last.message.channel, last.message.ts);
167187
}
188+
throw error;
168189
}
169-
await dispatchPreparedSlackMessage(prepared);
170190
},
171191
onError: (err) => {
172192
ctx.runtime.error?.(`slack inbound debounce flush failed: ${String(err)}`);

extensions/slack/src/monitor/message-handler/prepare.test.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -642,6 +642,7 @@ describe("prepareSlackMessage sender prefix", () => {
642642
removeAckAfterReply: false,
643643
logger: { info: vi.fn(), warn: vi.fn() },
644644
markMessageSeen: () => false,
645+
releaseSeenMessage: () => {},
645646
shouldDropMismatchedSlackEvent: () => false,
646647
resolveSlackSystemEventSessionKey: () => "agent:main:slack:channel:c1",
647648
isChannelAllowed: () => true,

0 commit comments

Comments
 (0)