From be7a2c23e97d4ebeeb8a9c59aebcbe5eb03d638f Mon Sep 17 00:00:00 2001 From: joshavant <830519+joshavant@users.noreply.github.com> Date: Thu, 13 Aug 2026 00:31:30 -0500 Subject: [PATCH] fix(discord): cancel started ingress jobs on abort --- .../src/monitor/message-handler.queue.test.ts | 90 +++++++++++++++++++ .../discord/src/monitor/message-run-queue.ts | 10 ++- 2 files changed, 98 insertions(+), 2 deletions(-) diff --git a/extensions/discord/src/monitor/message-handler.queue.test.ts b/extensions/discord/src/monitor/message-handler.queue.test.ts index 2b8bdea09ea8..ff8688fc70ff 100644 --- a/extensions/discord/src/monitor/message-handler.queue.test.ts +++ b/extensions/discord/src/monitor/message-handler.queue.test.ts @@ -639,6 +639,96 @@ describe("createDiscordMessageHandler queue behavior", () => { }); }); + it.each(["returns", "throws"] as const)( + "preserves retry facts when a started durable Discord job %s after cancellation", + async (outcome) => { + await withDiscordQueue(async (queue) => { + const id = `started-cancelled-${outcome}`; + const raw = createRawMessage(id, "lane-a"); + await queue.enqueue( + id, + { version: 1, receivedAt: 10, rawMessage: raw }, + { laneKey: "channel:lane-a", receivedAt: 10 }, + ); + const failedClaim = await queue.claim(id, { ownerId: "failed-owner" }); + expect(failedClaim).not.toBeNull(); + if (!failedClaim) { + return; + } + await queue.release(failedClaim, { + lastError: "previous genuine failure", + releasedAt: 20, + }); + const before = (await queue.listPending())[0]; + const processingStarted = createDeferred(); + const finishProcessing = createDeferred(); + let processingSignal: AbortSignal | undefined; + const processDiscordMessage = vi.fn(async (ctx: { abortSignal?: AbortSignal }) => { + processingSignal = ctx.abortSignal; + processingStarted.resolve(); + await finishProcessing.promise; + if (outcome === "throws") { + throw new Error("processing stopped after cancellation"); + } + }); + const params = createDiscordHandlerParams(); + const handler = createDurableDiscordMessageHandler({ + ...params, + client: {} as never, + testing: { + preflightDiscordMessage: (async (preflightParams: { + abortSignal?: AbortSignal; + data: ReturnType; + turnAdoptionLifecycle?: DiscordIngressLifecycle; + }) => ({ + ...createPreflightContextForMessage(preflightParams.data), + abortSignal: preflightParams.abortSignal, + turnAdoptionLifecycle: preflightParams.turnAdoptionLifecycle, + })) as never, + processDiscordMessage: processDiscordMessage as never, + createIngressMonitor: (monitorParams) => + createDiscordIngressMonitor({ ...monitorParams, queue }), + }, + }); + + await processingStarted.promise; + const deactivation = handler.deactivate(); + await vi.waitFor(() => expect(processingSignal?.aborted).toBe(true)); + finishProcessing.resolve(); + await deactivation; + + expect(await queue.listPending()).toEqual([ + expect.objectContaining({ + id, + attempts: before?.attempts, + lastAttemptAt: before?.lastAttemptAt, + lastError: before?.lastError, + }), + ]); + + const recovered = vi.fn(async (_event, lifecycle: DiscordIngressLifecycle) => { + await lifecycle.onAdopted(); + }); + const replacement = createDiscordIngressMonitor({ + accountId: "default", + client: {} as never, + runtime: params.runtime, + queue, + dispatch: recovered, + }); + replacement.start(); + try { + await vi.waitFor(() => expect(recovered).toHaveBeenCalledTimes(1)); + await expect(queue.enqueue(id, {} as DiscordIngressPayload)).resolves.toMatchObject({ + kind: "completed", + }); + } finally { + await replacement.stop(); + } + }); + }, + ); + it("preserves retry facts when deactivation skips a queued durable Discord job", async () => { await withDiscordQueue(async (queue) => { const raw = createRawMessage("queued-cancelled", "lane-a"); diff --git a/extensions/discord/src/monitor/message-run-queue.ts b/extensions/discord/src/monitor/message-run-queue.ts index f2752dd237b6..bb1be5365841 100644 --- a/extensions/discord/src/monitor/message-run-queue.ts +++ b/extensions/discord/src/monitor/message-run-queue.ts @@ -45,12 +45,18 @@ async function processDiscordQueuedMessage(params: { (await loadMessageProcessRuntime()).processDiscordMessage; await processDiscordMessageImpl(materializeDiscordInboundJob(params.job, abortSignal)); if (abortSignal?.aborted) { - await params.job.ingressSettlement?.abandon(abortSignal.reason); + // Cancellation ended ownership before delivery; retain prior retry facts + // so the durable claim can replay under a replacement lifecycle. + await params.job.ingressSettlement?.cancel(); } else { await params.job.ingressSettlement?.settle(); } } catch (error) { - await params.job.ingressSettlement?.abandon(error); + if (abortSignal?.aborted) { + await params.job.ingressSettlement?.cancel(); + } else { + await params.job.ingressSettlement?.abandon(error); + } throw error; } }