Skip to content

Commit 5fe0d46

Browse files
committed
fix(channels): make nack callbacks idempotent
1 parent e864155 commit 5fe0d46

2 files changed

Lines changed: 29 additions & 0 deletions

File tree

src/channels/message/lifecycle.test.ts

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -367,11 +367,36 @@ describe("message lifecycle primitives", () => {
367367

368368
const nackError = new Error("offset failed");
369369
await ctx.nack(nackError);
370+
await ctx.nack(new Error("duplicate failure"));
371+
expect(onNack).toHaveBeenCalledTimes(1);
370372
expect(onNack).toHaveBeenCalledWith(nackError);
371373
expect(ctx.ackState).toBe("nacked");
372374
expect(ctx.nackErrorMessage).toBe("offset failed");
373375
});
374376

377+
it("retries nack callbacks after a failed attempt", async () => {
378+
const onNack = vi
379+
.fn<(error: unknown) => Promise<void>>()
380+
.mockRejectedValueOnce(new Error("temporary failure"))
381+
.mockResolvedValueOnce();
382+
const ctx = createMessageReceiveContext({
383+
id: "rx-nack-retry",
384+
channel: "telegram",
385+
message: { text: "hello" },
386+
onNack,
387+
});
388+
389+
await expect(ctx.nack(new Error("first failure"))).rejects.toThrow("temporary failure");
390+
expect(ctx.ackState).toBe("pending");
391+
392+
const retryError = new Error("retry failure");
393+
await ctx.nack(retryError);
394+
395+
expect(onNack).toHaveBeenCalledTimes(2);
396+
expect(ctx.ackState).toBe("nacked");
397+
expect(ctx.nackErrorMessage).toBe("retry failure");
398+
});
399+
375400
it("maps ack policies to lifecycle stages", () => {
376401
expect(shouldAckMessageAfterStage("after_receive_record", "receive_record")).toBe(true);
377402
expect(shouldAckMessageAfterStage("after_receive_record", "agent_dispatch")).toBe(false);

src/channels/message/receive.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,10 @@ export function createMessageReceiveContext<TMessage>(params: {
8888
delete ctx.nackErrorMessage;
8989
},
9090
nack: async (error) => {
91+
// Nack callbacks are idempotent after completion but remain retryable after failure.
92+
if (ctx.ackState === "nacked") {
93+
return;
94+
}
9195
await params.onNack?.(error);
9296
ctx.ackState = "nacked";
9397
ctx.nackErrorMessage = normalizeAckErrorMessage(error);

0 commit comments

Comments
 (0)