From 6a96bbc7ad77bf60df4608d978d740c95e28eb65 Mon Sep 17 00:00:00 2001 From: Ayaan Zaidi Date: Wed, 1 Jul 2026 17:22:32 -0700 Subject: [PATCH] fix(telegram): keep timed-out webhook lanes guarded (#98806) --- extensions/telegram/src/webhook.test.ts | 54 ++++++++++++++++- extensions/telegram/src/webhook.ts | 80 ++++++++++++++++++------- 2 files changed, 113 insertions(+), 21 deletions(-) diff --git a/extensions/telegram/src/webhook.test.ts b/extensions/telegram/src/webhook.test.ts index e8a05ffd63cb..87e07ac27787 100644 --- a/extensions/telegram/src/webhook.test.ts +++ b/extensions/telegram/src/webhook.test.ts @@ -825,6 +825,55 @@ describe("startTelegramWebhook", () => { ); }); + it("keeps a timed-out webhook lane guarded until replay settles", async () => { + vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] }); + try { + let finishFirstUpdate: (() => void) | undefined; + const seenUpdateIds: number[] = []; + const firstUpdate = { update_id: 40, message: { chat: { id: 123 }, text: "slow" } }; + const secondUpdate = { update_id: 41, message: { chat: { id: 123 }, text: "blocked" } }; + await writeTelegramSpooledUpdate({ + spoolDir: requireWebhookSpoolDir(), + update: firstUpdate, + }); + await writeTelegramSpooledUpdate({ + spoolDir: requireWebhookSpoolDir(), + update: secondUpdate, + }); + handleUpdateSpy.mockImplementation(async (update: unknown) => { + const updateId = (update as { update_id: number }).update_id; + seenUpdateIds.push(updateId); + if (updateId === 40) { + await new Promise((resolve) => { + finishFirstUpdate = resolve; + }); + } + }); + + const started = await startTelegramWebhook({ + token: TELEGRAM_TOKEN, + port: 0, + secret: TELEGRAM_SECRET, + path: TELEGRAM_WEBHOOK_PATH, + spoolDir: requireWebhookSpoolDir(), + runtime: { log: vi.fn(), error: vi.fn(), exit: vi.fn() }, + }); + try { + await vi.waitFor(() => expect(seenUpdateIds).toEqual([40])); + await vi.advanceTimersByTimeAsync(25 * 60_000 + 10_000); + await yieldWebhookTask(); + expect(seenUpdateIds).toEqual([40]); + + finishFirstUpdate?.(); + await vi.waitFor(() => expect(seenUpdateIds).toEqual([40, 41])); + } finally { + await started.stop(); + } + } finally { + vi.useRealTimers(); + } + }); + it("drains spooled webhook updates left by a previous process on startup", async () => { const update = { update_id: 30, message: { text: "leftover" } }; await writeTelegramSpooledUpdate({ @@ -917,7 +966,10 @@ describe("startTelegramWebhook", () => { [], ), ); - expectMockMessageContains(runtimeLog, "reached retry limit after 8 attempts; dead-lettered"); + expectMockMessageContains( + runtimeLog, + "reached retry limit after 8 attempts; dead-lettered", + ); } finally { await started.stop(); } diff --git a/extensions/telegram/src/webhook.ts b/extensions/telegram/src/webhook.ts index f698fc427f66..0133551328ec 100644 --- a/extensions/telegram/src/webhook.ts +++ b/extensions/telegram/src/webhook.ts @@ -84,7 +84,11 @@ const TELEGRAM_WEBHOOK_REGISTRATION_RETRY_POLICY: BackoffPolicy = { factor: 2, jitter: 0.2, }; -const activeWebhookSpooledHandlersByLane = new Set(); +type ActiveWebhookSpooledHandler = { + laneKey: string; +}; + +const activeWebhookSpooledHandlersByLane = new Map(); function buildWebhookSpooledHandlerKey(params: { laneKey: string; spoolDir: string }): string { return `${params.spoolDir}\0${params.laneKey}`; @@ -93,9 +97,9 @@ function buildWebhookSpooledHandlerKey(params: { laneKey: string; spoolDir: stri function resolveActiveWebhookSpooledLaneKeys(spoolDir: string): Set { const laneKeys = new Set(); const prefix = `${spoolDir}\0`; - for (const handlerKey of activeWebhookSpooledHandlersByLane) { + for (const [handlerKey, handler] of activeWebhookSpooledHandlersByLane) { if (handlerKey.startsWith(prefix)) { - laneKeys.add(handlerKey.slice(prefix.length)); + laneKeys.add(handler.laneKey); } } return laneKeys; @@ -454,17 +458,21 @@ async function waitForTimedOutWebhookReplayGrace(params: { log: (line: string) => void; replayTask: Promise<{ deferredWork?: TelegramSpooledReplayDeferredParticipant }>; updateId: number; -}): Promise { +}): Promise { let timer: ReturnType | undefined; try { - await Promise.race([ - params.replayTask.catch((replayErr: unknown) => { - params.log( - `[telegram][diag] timed out webhook spooled update ${params.updateId} replay later failed: ${formatErrorMessage(replayErr)}`, - ); - }), - new Promise((resolve) => { - timer = setTimeout(resolve, TELEGRAM_WEBHOOK_SPOOLED_HANDLER_ABORT_GRACE_MS); + return await Promise.race([ + params.replayTask.then( + () => true, + (replayErr: unknown) => { + params.log( + `[telegram][diag] timed out webhook spooled update ${params.updateId} replay later failed: ${formatErrorMessage(replayErr)}`, + ); + return true; + }, + ), + new Promise((resolve) => { + timer = setTimeout(() => resolve(false), TELEGRAM_WEBHOOK_SPOOLED_HANDLER_ABORT_GRACE_MS); timer.unref?.(); }), ]); @@ -475,6 +483,10 @@ async function waitForTimedOutWebhookReplayGrace(params: { } } +type WebhookSpooledUpdateHandlerResult = { + retainLaneGuardTask?: Promise; +}; + async function runWebhookSpooledReplayWithTimeout(params: { bot: ReturnType; laneKey: string; @@ -549,7 +561,7 @@ async function handleWebhookSpooledUpdate(params: { bot: ReturnType; log: (line: string) => void; update: ClaimedTelegramSpooledUpdate; -}): Promise { +}): Promise { let replay: { deferredWork?: TelegramSpooledReplayDeferredParticipant }; try { const rawUpdate = params.update.update; @@ -583,19 +595,28 @@ async function handleWebhookSpooledUpdate(params: { message: err.message, update: params.update, }); - await waitForTimedOutWebhookReplayGrace({ + const replaySettled = await waitForTimedOutWebhookReplayGrace({ log: params.log, replayTask: err.replayTask, updateId: params.update.updateId, }); - return; + if (replaySettled) { + return {}; + } + return { + retainLaneGuardTask: err.replayTask.catch((replayErr: unknown) => { + params.log( + `[telegram][diag] timed out webhook spooled update ${params.update.updateId} replay later failed: ${formatErrorMessage(replayErr)}`, + ); + }), + }; } await releaseFailedWebhookSpooledUpdate({ err, log: params.log, update: params.update, }); - return; + return {}; } if (replay.deferredWork) { const result = await waitForWebhookSpooledDeferredWork({ @@ -611,14 +632,14 @@ async function handleWebhookSpooledUpdate(params: { message: formatErrorMessage(result.error), update: params.update, }); - return; + return {}; } await releaseFailedWebhookSpooledUpdate({ err: result.error, log: params.log, update: params.update, }); - return; + return {}; } } try { @@ -628,6 +649,7 @@ async function handleWebhookSpooledUpdate(params: { `[telegram][diag] webhook spooled update ${params.update.updateId} completed but processing marker cleanup failed: ${formatErrorMessage(err)}`, ); } + return {}; } export async function startTelegramWebhook(opts: { @@ -769,7 +791,8 @@ export async function startTelegramWebhook(opts: { const handlerKey = buildWebhookSpooledHandlerKey({ spoolDir, laneKey }); // Webhook HTTP requests and same-process restarts can overlap; keep // one process-global active claim per spool lane to preserve ordering. - activeWebhookSpooledHandlersByLane.add(handlerKey); + const handlerState: ActiveWebhookSpooledHandler = { laneKey }; + activeWebhookSpooledHandlersByLane.set(handlerKey, handlerState); blockedLaneKeys.add(laneKey); // Claim ownership has a finite lease; refresh while the handler runs so // another process cannot recover and replay this update concurrently. @@ -777,12 +800,24 @@ export async function startTelegramWebhook(opts: { log, update: claimedUpdate, }); + let retainLaneGuardTask: Promise | undefined; void handleWebhookSpooledUpdate({ accountId: opts.accountId ?? "default", bot, log, update: claimedUpdate, }) + .then((result) => { + retainLaneGuardTask = result.retainLaneGuardTask; + if (retainLaneGuardTask) { + void retainLaneGuardTask.finally(() => { + if (activeWebhookSpooledHandlersByLane.get(handlerKey) === handlerState) { + activeWebhookSpooledHandlersByLane.delete(handlerKey); + } + void Promise.resolve().then(drainWebhookSpool); + }); + } + }) .catch((err: unknown) => { runtime.log?.( `[telegram][diag] webhook spooled update ${claimedUpdate.updateId} handler failed after claim: ${formatErrorMessage(err)}`, @@ -790,7 +825,12 @@ export async function startTelegramWebhook(opts: { }) .finally(() => { stopClaimRefresh(); - activeWebhookSpooledHandlersByLane.delete(handlerKey); + if ( + !retainLaneGuardTask && + activeWebhookSpooledHandlersByLane.get(handlerKey) === handlerState + ) { + activeWebhookSpooledHandlersByLane.delete(handlerKey); + } void Promise.resolve().then(drainWebhookSpool); }); started += 1;