Skip to content

Commit 1440ed9

Browse files
committed
fix(discord): unblock inbound queue after stuck run
1 parent cec50aa commit 1440ed9

4 files changed

Lines changed: 262 additions & 17 deletions

File tree

extensions/discord/src/monitor/message-handler.queue.test.ts

Lines changed: 170 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,10 @@ import {
1111
createDiscordPreflightContext,
1212
} from "./message-handler.test-helpers.js";
1313

14+
const agentTimeoutMocks = vi.hoisted(() => ({
15+
resolveAgentTimeoutMs: vi.fn(() => 48 * 60 * 60 * 1000),
16+
}));
17+
1418
const earlyTypingMocks = vi.hoisted(() => ({
1519
createDiscordRestClient: vi.fn(() => ({
1620
token: "test-token",
@@ -20,6 +24,10 @@ const earlyTypingMocks = vi.hoisted(() => ({
2024
sendTyping: vi.fn(async () => {}),
2125
}));
2226

27+
vi.mock("openclaw/plugin-sdk/agent-runtime", () => ({
28+
resolveAgentTimeoutMs: agentTimeoutMocks.resolveAgentTimeoutMs,
29+
}));
30+
2331
vi.mock("../client.js", () => ({
2432
createDiscordRestClient: earlyTypingMocks.createDiscordRestClient,
2533
}));
@@ -172,6 +180,7 @@ async function createLifecycleStopScenario(params: {
172180

173181
describe("createDiscordMessageHandler queue behavior", () => {
174182
beforeEach(() => {
183+
agentTimeoutMocks.resolveAgentTimeoutMs.mockReset().mockReturnValue(48 * 60 * 60 * 1000);
175184
earlyTypingMocks.createDiscordRestClient.mockReset().mockReturnValue({
176185
token: "test-token",
177186
rest: { kind: "discord-rest" },
@@ -420,7 +429,7 @@ describe("createDiscordMessageHandler queue behavior", () => {
420429
expect(visibleSideEffect).toHaveBeenCalledTimes(1);
421430
});
422431

423-
it("does not abort long queued runs with a Discord-owned channel timeout", async () => {
432+
it("does not abort long queued runs before the Discord run watchdog expires", async () => {
424433
vi.useFakeTimers();
425434
try {
426435
preflightDiscordMessageMock.mockReset();
@@ -458,7 +467,8 @@ describe("createDiscordMessageHandler queue behavior", () => {
458467
await flushQueueWork();
459468

460469
expect(processDiscordMessageMock).toHaveBeenCalledTimes(1);
461-
expect(capturedAbortSignals).toEqual([undefined]);
470+
expect(capturedAbortSignals).toHaveLength(1);
471+
expect(capturedAbortSignals[0]?.aborted).toBe(false);
462472
const runtimeError = params.runtime.error as unknown as MockCallSource;
463473
expect(
464474
mockCalls(runtimeError).some(([message]) => String(message).includes("timed out")),
@@ -469,7 +479,164 @@ describe("createDiscordMessageHandler queue behavior", () => {
469479
await flushQueueWork();
470480

471481
expect(processDiscordMessageMock).toHaveBeenCalledTimes(2);
472-
expect(capturedAbortSignals).toEqual([undefined, undefined]);
482+
expect(capturedAbortSignals).toHaveLength(2);
483+
expect(capturedAbortSignals[1]?.aborted).toBe(false);
484+
485+
secondRun.resolve();
486+
await secondRun.promise;
487+
} finally {
488+
vi.useRealTimers();
489+
}
490+
});
491+
492+
it("uses the configured agent timeout as the default Discord run watchdog", async () => {
493+
vi.useFakeTimers();
494+
try {
495+
preflightDiscordMessageMock.mockReset();
496+
processDiscordMessageMock.mockReset();
497+
498+
const run = createDeferred();
499+
const ctx = createPreflightContext();
500+
preflightDiscordMessageMock.mockResolvedValue(ctx);
501+
processDiscordMessageMock.mockImplementation(async () => {
502+
await run.promise;
503+
});
504+
const handler = createDiscordMessageHandler(createDiscordHandlerParams());
505+
506+
await expect(
507+
handler(createMessageData("m-configured") as never, {} as never),
508+
).resolves.toBeUndefined();
509+
await flushQueueWork();
510+
511+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(1);
512+
expect(agentTimeoutMocks.resolveAgentTimeoutMs).toHaveBeenCalledWith({ cfg: ctx.cfg });
513+
514+
await vi.advanceTimersByTimeAsync(60_000);
515+
await flushQueueWork();
516+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(1);
517+
518+
run.resolve();
519+
await run.promise;
520+
} finally {
521+
vi.useRealTimers();
522+
}
523+
});
524+
525+
it("aborts a stuck Discord run and lets later messages for the same session proceed", async () => {
526+
vi.useFakeTimers();
527+
try {
528+
preflightDiscordMessageMock.mockReset();
529+
processDiscordMessageMock.mockReset();
530+
531+
const firstRun = createDeferred();
532+
const secondProcessed = vi.fn();
533+
const capturedAbortSignals: Array<AbortSignal | undefined> = [];
534+
processDiscordMessageMock.mockImplementationOnce(
535+
async (ctx: { abortSignal?: AbortSignal }) => {
536+
capturedAbortSignals.push(ctx.abortSignal);
537+
await firstRun.promise;
538+
},
539+
);
540+
processDiscordMessageMock.mockImplementationOnce(
541+
async (ctx: { abortSignal?: AbortSignal }) => {
542+
capturedAbortSignals.push(ctx.abortSignal);
543+
secondProcessed();
544+
},
545+
);
546+
installDefaultDiscordPreflight();
547+
const params = createDiscordHandlerParams();
548+
const handler = createDiscordMessageHandler({
549+
...params,
550+
testing: { runTimeoutMs: 25 },
551+
});
552+
553+
await expect(
554+
handler(createMessageData("m-1") as never, {} as never),
555+
).resolves.toBeUndefined();
556+
await expect(
557+
handler(createMessageData("m-2") as never, {} as never),
558+
).resolves.toBeUndefined();
559+
await flushQueueWork();
560+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(1);
561+
562+
await vi.advanceTimersByTimeAsync(25);
563+
await flushQueueWork();
564+
565+
expect(capturedAbortSignals[0]?.aborted).toBe(true);
566+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(2);
567+
expect(secondProcessed).toHaveBeenCalledTimes(1);
568+
const runtimeError = params.runtime.error as unknown as MockCallSource;
569+
expect(
570+
mockCalls(runtimeError).some(([message]) =>
571+
String(message).includes("DiscordMessageRunTimeoutError"),
572+
),
573+
).toBe(true);
574+
575+
await expect(
576+
handler(createMessageData("m-1") as never, {} as never),
577+
).resolves.toBeUndefined();
578+
await flushQueueWork();
579+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(2);
580+
expect(preflightDiscordMessageMock).toHaveBeenCalledTimes(2);
581+
582+
firstRun.resolve();
583+
await firstRun.promise;
584+
await flushQueueWork();
585+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(2);
586+
expect(params.runtime.error).toHaveBeenCalledTimes(1);
587+
} finally {
588+
vi.useRealTimers();
589+
}
590+
});
591+
592+
it("can disable the Discord run watchdog with a zero timeout", async () => {
593+
vi.useFakeTimers();
594+
try {
595+
preflightDiscordMessageMock.mockReset();
596+
processDiscordMessageMock.mockReset();
597+
598+
const firstRun = createDeferred();
599+
const secondRun = createDeferred();
600+
const capturedAbortSignals: Array<AbortSignal | undefined> = [];
601+
processDiscordMessageMock.mockImplementationOnce(
602+
async (ctx: { abortSignal?: AbortSignal }) => {
603+
capturedAbortSignals.push(ctx.abortSignal);
604+
await firstRun.promise;
605+
},
606+
);
607+
processDiscordMessageMock.mockImplementationOnce(
608+
async (ctx: { abortSignal?: AbortSignal }) => {
609+
capturedAbortSignals.push(ctx.abortSignal);
610+
await secondRun.promise;
611+
},
612+
);
613+
installDefaultDiscordPreflight();
614+
const params = createDiscordHandlerParams();
615+
const handler = createDiscordMessageHandler({
616+
...params,
617+
testing: { runTimeoutMs: 0 },
618+
});
619+
620+
await expect(
621+
handler(createMessageData("m-1") as never, {} as never),
622+
).resolves.toBeUndefined();
623+
await expect(
624+
handler(createMessageData("m-2") as never, {} as never),
625+
).resolves.toBeUndefined();
626+
await flushQueueWork();
627+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(1);
628+
629+
await vi.advanceTimersByTimeAsync(60 * 60 * 1000);
630+
await flushQueueWork();
631+
632+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(1);
633+
expect(capturedAbortSignals).toEqual([undefined]);
634+
expect(params.runtime.error).not.toHaveBeenCalled();
635+
636+
firstRun.resolve();
637+
await firstRun.promise;
638+
await flushQueueWork();
639+
expect(processDiscordMessageMock).toHaveBeenCalledTimes(2);
473640

474641
secondRun.resolve();
475642
await secondRun.promise;

extensions/discord/src/monitor/message-handler.test-helpers.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,10 @@ export function createDiscordHandlerParams(overrides?: {
5353

5454
export function createDiscordPreflightContext(channelId = "ch-1") {
5555
return {
56+
runtime: {
57+
log: vi.fn(),
58+
error: vi.fn(),
59+
},
5660
data: {
5761
channel_id: channelId,
5862
message: {

extensions/discord/src/monitor/message-run-queue.ts

Lines changed: 80 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import { resolveAgentTimeoutMs } from "openclaw/plugin-sdk/agent-runtime";
12
import { createChannelRunQueue } from "openclaw/plugin-sdk/channel-outbound";
23
import type { ClaimableDedupe } from "openclaw/plugin-sdk/persistent-dedupe";
34
import { danger } from "openclaw/plugin-sdk/runtime-env";
@@ -14,6 +15,13 @@ import { mergeAbortSignals } from "./timeouts.js";
1415

1516
type ProcessDiscordMessage = typeof import("./message-handler.process.js").processDiscordMessage;
1617

18+
class DiscordMessageRunTimeoutError extends Error {
19+
constructor(timeoutMs: number) {
20+
super(`Discord message run timed out after ${timeoutMs}ms`);
21+
this.name = "DiscordMessageRunTimeoutError";
22+
}
23+
}
24+
1725
type DiscordMessageRunQueueParams = {
1826
runtime: RuntimeEnv;
1927
setStatus?: DiscordMonitorStatusSink;
@@ -29,6 +37,7 @@ type DiscordMessageRunQueue = {
2937

3038
export type DiscordMessageRunQueueTestingHooks = {
3139
processDiscordMessage?: ProcessDiscordMessage;
40+
runTimeoutMs?: number;
3241
};
3342

3443
let messageProcessRuntimePromise:
@@ -40,6 +49,16 @@ async function loadMessageProcessRuntime() {
4049
return await messageProcessRuntimePromise;
4150
}
4251

52+
function normalizeDiscordMessageRunTimeoutMs(params: {
53+
value: unknown;
54+
job: DiscordInboundJob;
55+
}): number {
56+
if (typeof params.value === "number" && Number.isFinite(params.value)) {
57+
return Math.max(0, Math.floor(params.value));
58+
}
59+
return resolveAgentTimeoutMs({ cfg: params.job.payload.cfg });
60+
}
61+
4362
async function processDiscordQueuedMessage(params: {
4463
job: DiscordInboundJob;
4564
lifecycleSignal?: AbortSignal;
@@ -49,27 +68,76 @@ async function processDiscordQueuedMessage(params: {
4968
const processDiscordMessageImpl =
5069
params.testing?.processDiscordMessage ??
5170
(await loadMessageProcessRuntime()).processDiscordMessage;
52-
const abortSignal = mergeAbortSignals([params.job.runtime.abortSignal, params.lifecycleSignal]);
53-
try {
54-
await processDiscordMessageImpl(materializeDiscordInboundJob(params.job, abortSignal));
55-
await commitDiscordInboundReplay({
56-
replayKeys: params.job.replayKeys,
57-
replayGuard: params.replayGuard,
58-
});
59-
} catch (error) {
60-
if (error instanceof DiscordRetryableInboundError) {
61-
releaseDiscordInboundReplay({
71+
const runTimeoutMs = normalizeDiscordMessageRunTimeoutMs({
72+
value: params.testing?.runTimeoutMs,
73+
job: params.job,
74+
});
75+
const timeoutController = runTimeoutMs > 0 ? new AbortController() : undefined;
76+
const abortSignal = mergeAbortSignals([
77+
params.job.runtime.abortSignal,
78+
params.lifecycleSignal,
79+
timeoutController?.signal,
80+
]);
81+
let timeoutId: ReturnType<typeof setTimeout> | undefined;
82+
let timedOut = false;
83+
const processPromise = (async () => {
84+
try {
85+
await processDiscordMessageImpl(materializeDiscordInboundJob(params.job, abortSignal));
86+
if (timedOut) {
87+
return;
88+
}
89+
await commitDiscordInboundReplay({
6290
replayKeys: params.job.replayKeys,
63-
error,
6491
replayGuard: params.replayGuard,
6592
});
66-
} else {
93+
} catch (error) {
94+
if (timedOut) {
95+
return;
96+
}
97+
if (error instanceof DiscordRetryableInboundError) {
98+
releaseDiscordInboundReplay({
99+
replayKeys: params.job.replayKeys,
100+
error,
101+
replayGuard: params.replayGuard,
102+
});
103+
} else {
104+
await commitDiscordInboundReplay({
105+
replayKeys: params.job.replayKeys,
106+
replayGuard: params.replayGuard,
107+
});
108+
}
109+
throw error;
110+
}
111+
})();
112+
try {
113+
if (runTimeoutMs <= 0) {
114+
await processPromise;
115+
return;
116+
}
117+
await Promise.race([
118+
processPromise,
119+
new Promise<never>((_resolve, reject) => {
120+
timeoutId = setTimeout(() => {
121+
timedOut = true;
122+
const error = new DiscordMessageRunTimeoutError(runTimeoutMs);
123+
timeoutController?.abort(error);
124+
reject(error);
125+
}, runTimeoutMs);
126+
timeoutId.unref?.();
127+
}),
128+
]);
129+
} catch (error) {
130+
if (error instanceof DiscordMessageRunTimeoutError) {
67131
await commitDiscordInboundReplay({
68132
replayKeys: params.job.replayKeys,
69133
replayGuard: params.replayGuard,
70134
});
71135
}
72136
throw error;
137+
} finally {
138+
if (timeoutId !== undefined) {
139+
clearTimeout(timeoutId);
140+
}
73141
}
74142
}
75143

extensions/discord/src/monitor/timeouts.ts

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,14 @@
11
import { resolveTimerTimeoutMs } from "openclaw/plugin-sdk/number-runtime";
22

3-
// Compatibility constants for existing imports. Discord no longer enforces
4-
// channel-owned listener or inbound run timeouts.
3+
// Compatibility constant for existing imports. Discord no longer uses this
4+
// listener timeout to cover full agent/message processing.
55
export const DISCORD_DEFAULT_LISTENER_TIMEOUT_MS = 120_000;
6+
7+
/**
8+
* @deprecated Compatibility export for callers that imported the previous
9+
* Discord inbound worker default. Discord message-run watchdogs now use the
10+
* configured agent timeout resolver instead of this constant.
11+
*/
612
export const DISCORD_DEFAULT_INBOUND_WORKER_TIMEOUT_MS = 30 * 60_000;
713

814
export const DISCORD_ATTACHMENT_IDLE_TIMEOUT_MS = 60_000;

0 commit comments

Comments
 (0)