From e29381a172ec135aaa7db54b0ecd09ef55a538b1 Mon Sep 17 00:00:00 2001 From: weiqinl Date: Sun, 21 Jun 2026 23:10:45 +0800 Subject: [PATCH] fix #86957: drain worker-spooled Telegram updates immediately Wake the isolated polling drain immediately after a worker-spooled update is written to channel_ingress_events, instead of waiting for the next drain interval. - Add requestImmediateDrain() calls after worker write and spooled message - Track pending drain requests while drain is active (fix race condition) - Add regression test for updates arriving during active drain Fixes #86957. --- .../telegram/src/polling-session.test.ts | 166 ++++++++++++++++++ extensions/telegram/src/polling-session.ts | 18 +- 2 files changed, 183 insertions(+), 1 deletion(-) diff --git a/extensions/telegram/src/polling-session.test.ts b/extensions/telegram/src/polling-session.test.ts index 51c7f75baad7..6fe8cab93974 100644 --- a/extensions/telegram/src/polling-session.test.ts +++ b/extensions/telegram/src/polling-session.test.ts @@ -903,6 +903,172 @@ describe("TelegramPollingSession", () => { } }); + it("drains worker-spooled updates without waiting for the next drain interval", async () => { + const abort = new AbortController(); + const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-")); + const handleUpdate = vi.fn(async () => { + abort.abort(); + }); + const bot = { + api: { + deleteWebhook: vi.fn(async () => true), + config: { use: vi.fn() }, + }, + init: vi.fn(async () => undefined), + handleUpdate, + stop: vi.fn(async () => undefined), + }; + createTelegramBotMock.mockReturnValueOnce(bot); + let onMessage: WorkerMessageListener | undefined; + let stopWorker: (() => void) | undefined; + const workerDone = new Promise((resolve) => { + stopWorker = resolve; + }); + const ackSpooledUpdate = vi.fn(); + const createWorker = vi.fn(() => ({ + onMessage: vi.fn((listener: WorkerMessageListener) => { + onMessage = listener; + return () => undefined; + }), + ackSpooledUpdate, + stop: vi.fn(async () => { + stopWorker?.(); + }), + task: vi.fn(async () => { + await workerDone; + }), + })); + + try { + const session = createPollingSession({ + abortSignal: abort.signal, + isolatedIngress: { + enabled: true, + spoolDir: tempDir, + createWorker, + drainIntervalMs: 60_000, + }, + }); + + const runPromise = session.runUntilAbort(); + await vi.waitFor(() => expect(onMessage).toBeDefined()); + onMessage?.({ + type: "update", + requestId: "write-1", + update: { update_id: 42, message: { text: "hello" } }, + queued: 1, + }); + + await vi.waitFor(() => + expect(ackSpooledUpdate).toHaveBeenCalledWith("write-1", { ok: true, updateId: 42 }), + ); + onMessage?.({ type: "spooled", updateId: 42, queued: 1 }); + await vi.waitFor(() => + expect(handleUpdate).toHaveBeenCalledWith({ update_id: 42, message: { text: "hello" } }), + ); + await vi.waitFor(async () => expect(await pendingUpdateIds(tempDir, "all")).toEqual([])); + stopWorker?.(); + await runPromise; + } finally { + abort.abort(); + stopWorker?.(); + await fs.rm(tempDir, { recursive: true, force: true }); + } + }); + + it("drains worker-spooled updates that arrive during an active drain", async () => { + const abort = new AbortController(); + const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-")); + + await writeTelegramSpooledUpdate({ + spoolDir: tempDir, + update: { update_id: 1, message: { text: "pre-seeded" } }, + }); + + const handleUpdate = vi.fn(async (update) => { + if (update.update_id === 1) { + await new Promise((resolve) => { setTimeout(resolve, 300); }); + } + }); + + const bot = { + api: { + deleteWebhook: vi.fn(async () => true), + config: { use: vi.fn() }, + }, + init: vi.fn(async () => undefined), + handleUpdate, + stop: vi.fn(async () => undefined), + }; + createTelegramBotMock.mockReturnValueOnce(bot); + let onMessage: WorkerMessageListener | undefined; + let stopWorker: (() => void) | undefined; + const workerDone = new Promise((resolve) => { + stopWorker = resolve; + }); + const ackSpooledUpdate = vi.fn(); + const createWorker = vi.fn(() => ({ + onMessage: vi.fn((listener: WorkerMessageListener) => { + onMessage = listener; + return () => undefined; + }), + ackSpooledUpdate, + stop: vi.fn(async () => { + stopWorker?.(); + }), + task: vi.fn(async () => { + await workerDone; + }), + })); + + try { + const session = createPollingSession({ + abortSignal: abort.signal, + isolatedIngress: { + enabled: true, + spoolDir: tempDir, + createWorker, + drainIntervalMs: 60_000, + }, + }); + + const runPromise = session.runUntilAbort(); + + await vi.waitFor(() => + expect(handleUpdate).toHaveBeenCalledWith({ + update_id: 1, + message: { text: "pre-seeded" }, + }), + ); + + onMessage?.({ + type: "update", + requestId: "write-2", + update: { update_id: 2, message: { text: "during-drain" } }, + queued: 1, + }); + + await vi.waitFor(() => + expect(ackSpooledUpdate).toHaveBeenCalledWith("write-2", { ok: true, updateId: 2 }), + ); + onMessage?.({ type: "spooled", updateId: 2, queued: 1 }); + + await vi.waitFor(() => + expect(handleUpdate).toHaveBeenCalledWith({ + update_id: 2, + message: { text: "during-drain" }, + }), + ); + await vi.waitFor(async () => expect(await pendingUpdateIds(tempDir, "all")).toEqual([])); + stopWorker?.(); + await runPromise; + } finally { + abort.abort(); + stopWorker?.(); + await fs.rm(tempDir, { recursive: true, force: true }); + } + }); + it("drains existing isolated ingress spool entries below the persisted offset", async () => { const abort = new AbortController(); const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-")); diff --git a/extensions/telegram/src/polling-session.ts b/extensions/telegram/src/polling-session.ts index bad94be8fa47..5a4f72b44395 100644 --- a/extensions/telegram/src/polling-session.ts +++ b/extensions/telegram/src/polling-session.ts @@ -1000,6 +1000,8 @@ export class TelegramPollingSession { forceCycleResolve = resolve; }); const stalledBacklogKeys = new Set(); + let requestImmediateDrain: () => void = () => undefined; + let drainRequested = false; const unsubscribe = worker.onMessage((message) => { const ackSpooledUpdate = ( requestId: string, @@ -1053,6 +1055,7 @@ export class TelegramPollingSession { }).then( (updateId) => { ackSpooledUpdate(message.requestId, { ok: true, updateId }); + requestImmediateDrain(); }, (err: unknown) => { ackSpooledUpdate(message.requestId, { @@ -1065,6 +1068,7 @@ export class TelegramPollingSession { } if (message.type === "spooled") { liveness.noteGetUpdatesActivity(); + requestImmediateDrain(); } }); const stopOnAbort = () => { @@ -1105,10 +1109,15 @@ export class TelegramPollingSession { } }; const drainOnce = async () => { - if (restartRequested || drainActive || this.opts.abortSignal?.aborted) { + if (restartRequested || this.opts.abortSignal?.aborted) { + return; + } + if (drainActive) { + drainRequested = true; return; } drainActive = true; + drainRequested = false; try { const drain = await this.#drainSpooledUpdates({ bot, spoolDir }); consecutiveDrainFailures = 0; @@ -1151,8 +1160,15 @@ export class TelegramPollingSession { ); } finally { drainActive = false; + if (drainRequested && !restartRequested && !this.opts.abortSignal?.aborted) { + drainRequested = false; + void drainOnce(); + } } }; + requestImmediateDrain = () => { + void drainOnce(); + }; await drainOnce(); const drainTimer = setInterval(() => { void drainOnce();