From 4370233f0c99793b2e61d2000dcdb55eee24dacc Mon Sep 17 00:00:00 2001 From: Vincent Koc <25068+vincentkoc@users.noreply.github.com> Date: Fri, 10 Jul 2026 05:22:26 -0700 Subject: [PATCH] fix(telegram): serialize spool timeout with adoption (cherry picked from commit bdbac2a7e3b2fe1322a31dc4708167c4da800cb4) --- .../telegram/src/bot-handlers.runtime.ts | 54 ++++++++++----- extensions/telegram/src/bot-message.test.ts | 5 +- .../telegram/src/bot-processing-outcome.ts | 47 +++++++++++-- .../src/bot.create-telegram-bot.test.ts | 15 ++-- .../telegram/src/message-dispatch-dedupe.ts | 14 ++-- .../telegram/src/polling-session.test.ts | 69 +++++++++++++++++++ extensions/telegram/src/polling-session.ts | 1 - extensions/telegram/src/webhook.test.ts | 58 ++++++++++++++++ extensions/telegram/src/webhook.ts | 47 +++++++------ 9 files changed, 252 insertions(+), 58 deletions(-) diff --git a/extensions/telegram/src/bot-handlers.runtime.ts b/extensions/telegram/src/bot-handlers.runtime.ts index 4895eab3d8d5..62e00fe3963b 100644 --- a/extensions/telegram/src/bot-handlers.runtime.ts +++ b/extensions/telegram/src/bot-handlers.runtime.ts @@ -93,6 +93,7 @@ import { recordTelegramMessageProcessingResult, type TelegramMessageProcessingResult, type TelegramSpooledReplayDeferredParticipant, + type TelegramSpooledReplaySettlementHold, } from "./bot-processing-outcome.js"; import { MEDIA_GROUP_TIMEOUT_MS, @@ -405,6 +406,31 @@ export const registerTelegramHandlers = ({ participant.settle(result); } }; + const beginSpooledReplaySettlementHolds = ( + participants: readonly TelegramSpooledReplayDeferredParticipant[], + ) => { + const holds: TelegramSpooledReplaySettlementHold[] = []; + for (const participant of new Set(participants)) { + const hold = participant.beginSettlementHold(); + if (!hold) { + for (const acquired of holds) { + acquired.release("replay-pending"); + } + const reason = participant.abortSignal.reason; + throw reason instanceof Error + ? reason + : new Error( + `telegram spooled replay participant ${participant.key} settled before durable adoption`, + ); + } + holds.push(hold); + } + return (mode: Parameters[0]) => { + for (const hold of holds) { + hold.release(mode); + } + }; + }; const createSpooledReplayParticipantForBufferedWork = ( key: string, ): TelegramSpooledReplayDeferredParticipant | undefined => @@ -1473,8 +1499,6 @@ export const registerTelegramHandlers = ({ let dispatchDedupeCommitted = false; let spooledReplayFinalResult: TelegramMessageProcessingResult | undefined; let spooledReplayFinalization: Promise | undefined; - let spooledReplayAdoptionCommitInFlight = false; - let deferredProcessingCancellation: TelegramMessageProcessingResult | undefined; const spooledReplay = params.options?.spooledReplay === true || isTelegramSpooledReplayUpdate(params.ctx.update) || @@ -1490,6 +1514,10 @@ export const registerTelegramHandlers = ({ ) ?? undefined) : undefined; + const ingressSpooledReplayParticipants = [ + ...explicitParticipants, + ...(frameParticipant ? [frameParticipant] : []), + ]; const processingParticipant = explicitParticipants.length > 0 ? createTelegramSpooledReplayParticipant( @@ -1499,18 +1527,13 @@ export const registerTelegramHandlers = ({ if (processingParticipant && explicitParticipants.length > 0) { for (const participant of explicitParticipants) { void participant.task.then((result) => { - if (spooledReplayAdoptionCommitInFlight && result.kind !== "completed") { - deferredProcessingCancellation ??= result; - return; - } processingParticipant.settle(result); }); } } const spooledReplayParticipants = [ ...new Set([ - ...explicitParticipants, - ...(frameParticipant ? [frameParticipant] : []), + ...ingressSpooledReplayParticipants, ...(processingParticipant ? [processingParticipant] : []), ]), ]; @@ -1528,20 +1551,18 @@ export const registerTelegramHandlers = ({ if (result.kind === "completed") { // Do not cache or settle a durable-adoption failure. Deferred queue // ownership retries this callback with the same spool participants. - spooledReplayAdoptionCommitInFlight = true; + const releaseSettlementHolds = beginSpooledReplaySettlementHolds( + ingressSpooledReplayParticipants, + ); try { await commitDispatchDedupeKeys(params.dispatchDedupeKeys ?? [], { requirePersistent: true, }); } catch (error) { - spooledReplayAdoptionCommitInFlight = false; - if (deferredProcessingCancellation) { - processingParticipant?.settle(deferredProcessingCancellation); - } + releaseSettlementHolds("replay-pending"); throw error; } - spooledReplayAdoptionCommitInFlight = false; - deferredProcessingCancellation = undefined; + releaseSettlementHolds("discard-pending"); dispatchDedupeCommitted = true; } else { releaseDispatchDedupeKeys( @@ -1671,7 +1692,8 @@ export const registerTelegramHandlers = ({ }, spooledReplayAbortSignal: params.spooledReplayAbortSignal, spooledReplayParticipant: processingParticipant, - finalizeSpooledReplayResult: async (result) => await finalizeSpooledReplayResult(result), + finalizeSpooledReplayResult: async (processingResult) => + await finalizeSpooledReplayResult(processingResult), completeSpooledReplayAfterIrrevocableAdoption: async () => { const completed = { kind: "completed" } satisfies TelegramMessageProcessingResult; return await finalizeSpooledReplayResult(completed); diff --git a/extensions/telegram/src/bot-message.test.ts b/extensions/telegram/src/bot-message.test.ts index e0432552d030..bb7f308fe892 100644 --- a/extensions/telegram/src/bot-message.test.ts +++ b/extensions/telegram/src/bot-message.test.ts @@ -1,10 +1,7 @@ // Telegram tests cover bot message plugin behavior. import { beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; import type { TelegramBotDeps } from "./bot-deps.js"; -import type { - TelegramMessageProcessingResult, - TelegramSpooledReplayDeferredParticipant, -} from "./bot-processing-outcome.js"; +import type { TelegramMessageProcessingResult } from "./bot-processing-outcome.js"; const buildTelegramMessageContext = vi.hoisted(() => vi.fn()); const dispatchTelegramMessage = vi.hoisted(() => vi.fn()); diff --git a/extensions/telegram/src/bot-processing-outcome.ts b/extensions/telegram/src/bot-processing-outcome.ts index cf91f901ed7b..61251e47fd45 100644 --- a/extensions/telegram/src/bot-processing-outcome.ts +++ b/extensions/telegram/src/bot-processing-outcome.ts @@ -18,9 +18,15 @@ export type TelegramSpooledReplayDeferredParticipant = { key: string; abortSignal: AbortSignal; task: Promise; + /** Defers external timeout settlement while durable adoption decides ownership. */ + beginSettlementHold: () => TelegramSpooledReplaySettlementHold | undefined; settle: (result: TelegramMessageProcessingResult) => void; }; +export type TelegramSpooledReplaySettlementHold = { + release: (mode: "discard-pending" | "replay-pending") => void; +}; + const telegramUpdateProcessingFrames = new AsyncLocalStorage(); const telegramSpooledReplayFrames = new AsyncLocalStorage(); const telegramSpooledReplayUpdates = new WeakSet(); @@ -64,23 +70,56 @@ export function createTelegramSpooledReplayParticipant( ): TelegramSpooledReplayDeferredParticipant { const abortController = new AbortController(); let settled = false; + let settlementHeld = false; + let pendingSettlement: TelegramMessageProcessingResult | undefined; let resolveTask: (result: TelegramMessageProcessingResult) => void = () => {}; const task = new Promise((resolve) => { resolveTask = resolve; }); + const settleNow = (result: TelegramMessageProcessingResult) => { + if (settled) { + return; + } + settled = true; + if (result.kind !== "completed") { + abortController.abort(result.kind === "failed-retryable" ? result.error : result.kind); + } + resolveTask(result); + }; return { key, abortSignal: abortController.signal, task, + beginSettlementHold: () => { + if (settled || settlementHeld) { + return undefined; + } + settlementHeld = true; + let released = false; + return { + release: (mode) => { + if (released) { + return; + } + released = true; + settlementHeld = false; + const pending = pendingSettlement; + pendingSettlement = undefined; + if (mode === "replay-pending" && pending) { + settleNow(pending); + } + }, + }; + }, settle: (result) => { if (settled) { return; } - settled = true; - if (result.kind !== "completed") { - abortController.abort(result.kind === "failed-retryable" ? result.error : result.kind); + if (settlementHeld) { + pendingSettlement ??= result; + return; } - resolveTask(result); + settleNow(result); }, }; } diff --git a/extensions/telegram/src/bot.create-telegram-bot.test.ts b/extensions/telegram/src/bot.create-telegram-bot.test.ts index 6566d990559a..5e4b8a419a09 100644 --- a/extensions/telegram/src/bot.create-telegram-bot.test.ts +++ b/extensions/telegram/src/bot.create-telegram-bot.test.ts @@ -1360,19 +1360,24 @@ describe("createTelegramBot", () => { const queuedTurn = runQueuedTurn?.(); await commitStarted; const timeoutError = new Error("spooled replay timed out during durable adoption"); - firstParticipant.settle({ kind: "failed-retryable", error: timeoutError }); - await expect(firstParticipant.task).resolves.toEqual({ - kind: "failed-retryable", - error: timeoutError, + let firstParticipantSettled = false; + void firstParticipant.task.then(() => { + firstParticipantSettled = true; }); + firstParticipant.settle({ kind: "failed-retryable", error: timeoutError }); await flushTelegramTestMicrotasks(); + expect(firstParticipantSettled).toBe(false); + expect(firstParticipant.abortSignal.aborted).toBe(false); expect(releaseSpy).not.toHaveBeenCalled(); releaseCommit?.(); await queuedTurn; expect(modelTurnRan).toBe(true); expect(queuedAbortSignal?.aborted).toBe(false); - await expect(secondParticipant.task).resolves.toEqual({ kind: "completed" }); + await expect(Promise.all([firstParticipant.task, secondParticipant.task])).resolves.toEqual([ + { kind: "completed" }, + { kind: "completed" }, + ]); expect(commitSpy).toHaveBeenCalledTimes(1); expect(releaseSpy).not.toHaveBeenCalled(); } finally { diff --git a/extensions/telegram/src/message-dispatch-dedupe.ts b/extensions/telegram/src/message-dispatch-dedupe.ts index 8e2f9d6aa384..1676e3faa3e0 100644 --- a/extensions/telegram/src/message-dispatch-dedupe.ts +++ b/extensions/telegram/src/message-dispatch-dedupe.ts @@ -1,6 +1,7 @@ // Telegram plugin module implements message dispatch dedupe behavior. import path from "node:path"; import type { Message } from "grammy/types"; +import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime"; import { createClaimableDedupe, type ClaimableDedupe } from "openclaw/plugin-sdk/persistent-dedupe"; import { normalizeStringEntries, uniqueStrings } from "openclaw/plugin-sdk/string-coerce-runtime"; @@ -158,9 +159,8 @@ export async function commitTelegramMessageDispatchReplay(params: { // can race rollback and recreate a key after it was forgotten. for (const [index, key] of keys.entries()) { let diskError: unknown; - let recorded = false; try { - recorded = await params.guard.commit( + const recorded = await params.guard.commit( key, params.requirePersistent === true ? { @@ -172,7 +172,12 @@ export async function commitTelegramMessageDispatchReplay(params: { : { namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE }, ); if (params.requirePersistent === true && diskError !== undefined) { - throw diskError; + throw diskError instanceof Error + ? diskError + : new Error(formatErrorMessage(diskError), { cause: diskError }); + } + if (recorded) { + committedKeys.push(key); } } catch (error) { for (const pendingKey of keys.slice(index + 1)) { @@ -215,9 +220,6 @@ export async function commitTelegramMessageDispatchReplay(params: { } throw error; } - if (recorded) { - committedKeys.push(key); - } } } diff --git a/extensions/telegram/src/polling-session.test.ts b/extensions/telegram/src/polling-session.test.ts index 12e806e216fe..7bfe0e31d9ab 100644 --- a/extensions/telegram/src/polling-session.test.ts +++ b/extensions/telegram/src/polling-session.test.ts @@ -80,6 +80,8 @@ type TelegramMessageProcessingResult = import("./bot-processing-outcome.js").TelegramMessageProcessingResult; type TelegramSpooledReplayDeferredParticipant = import("./bot-processing-outcome.js").TelegramSpooledReplayDeferredParticipant; +type TelegramSpooledReplaySettlementHold = + import("./bot-processing-outcome.js").TelegramSpooledReplaySettlementHold; let beginTelegramReplyFence: typeof import("./telegram-reply-fence.js").beginTelegramReplyFence; let buildTelegramReplyFenceLaneKey: typeof import("./telegram-reply-fence.js").buildTelegramReplyFenceLaneKey; let endTelegramReplyFence: typeof import("./telegram-reply-fence.js").endTelegramReplyFence; @@ -2213,6 +2215,73 @@ describe("TelegramPollingSession", () => { }); }); + it("keeps refreshing a buffered claim while timeout settlement waits for adoption", async () => { + const refreshHarness = installSpooledClaimRefreshHarness(); + await withTempSpool(async (tempDir) => { + const abort = new AbortController(); + const log = vi.fn(); + let participant: TelegramSpooledReplayDeferredParticipant | undefined; + let settlementHold: TelegramSpooledReplaySettlementHold | undefined; + await writeSpooledTestUpdates(tempDir, [topicUpdate(42, 10, "held adoption")]); + + const { runPromise, stopWorker } = startIsolatedIngressSession({ + abort, + spoolDir: tempDir, + log, + drainIntervalMs: 10, + spooledUpdateHandlerTimeoutMs: 20, + handleUpdate: async (update) => { + const createdParticipant = createTelegramSpooledReplayDeferredParticipant( + `test-held-adoption:${update.update_id}`, + ); + if (!createdParticipant) { + throw new Error("expected spooled replay participant"); + } + participant = createdParticipant; + settlementHold = createdParticipant.beginSettlementHold(); + if (!settlementHold) { + throw new Error("expected spooled replay settlement hold"); + } + }, + }); + + try { + await vi.waitFor(() => expect(participant).toBeDefined()); + const before = await claimedAtForUpdate(tempDir, 42); + await new Promise((resolve) => { + setTimeout(resolve, 50); + }); + + expect(participant?.abortSignal.aborted).toBe(false); + expect( + (await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).map( + (claim) => claim.updateId, + ), + ).toEqual([42]); + + refreshHarness.triggerRefresh(); + await vi.waitFor(async () => + expect(await claimedAtForUpdate(tempDir, 42)).toBeGreaterThan(before), + ); + + settlementHold?.release("discard-pending"); + participant?.settle({ kind: "completed" }); + await vi.waitFor(async () => + expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]), + ); + expect(await failedUpdateIds(tempDir)).toEqual([]); + expectLogExcludes(log, "pre-adoption timed out behind update 42"); + } finally { + settlementHold?.release("replay-pending"); + participant?.settle({ kind: "skipped" }); + abort.abort(); + stopWorker(); + refreshHarness.restore(); + await runPromise; + } + }); + }); + it("completes spooled row at adoption while a long turn is still settling (healthy long turn)", async () => { await withTempSpool(async (tempDir) => { const abort = new AbortController(); diff --git a/extensions/telegram/src/polling-session.ts b/extensions/telegram/src/polling-session.ts index b1e826f8e4c7..2051d1845e54 100644 --- a/extensions/telegram/src/polling-session.ts +++ b/extensions/telegram/src/polling-session.ts @@ -757,7 +757,6 @@ export class TelegramPollingSession { // Pre-adoption only: once the deferred participant settles at adoption, // this timer is cleared. A fire means ingress never adopted the turn. state.timedOutMessage = `Telegram isolated polling spool pre-adoption timed out behind update ${params.update.updateId} on lane ${params.laneKey} after ${age}; marking the update failed (handler-timeout) and keeping the claim out of retry.`; - state.stopClaimRefresh(); params.deferredWork.settle({ kind: "failed-retryable", error: new Error(state.timedOutMessage), diff --git a/extensions/telegram/src/webhook.test.ts b/extensions/telegram/src/webhook.test.ts index 38af02712812..8fa880ab3b21 100644 --- a/extensions/telegram/src/webhook.test.ts +++ b/extensions/telegram/src/webhook.test.ts @@ -12,6 +12,11 @@ import { } from "openclaw/plugin-sdk/plugin-state-test-runtime"; import { WEBHOOK_RATE_LIMIT_DEFAULTS } from "openclaw/plugin-sdk/webhook-ingress"; import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { + createTelegramSpooledReplayDeferredParticipant, + type TelegramSpooledReplayDeferredParticipant, + type TelegramSpooledReplaySettlementHold, +} from "./bot-processing-outcome.js"; import { clearTelegramRuntime, setTelegramRuntime } from "./runtime.js"; import type { TelegramRuntime } from "./runtime.types.js"; import { TELEGRAM_SPOOLED_RETRY_DEAD_LETTER_MIN_AGE_MS } from "./spooled-update-retry-policy.js"; @@ -928,6 +933,59 @@ describe("startTelegramWebhook", () => { } }); + it("holds buffered timeout settlement behind durable webhook adoption", async () => { + vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] }); + try { + const update = { update_id: 42, message: { chat: { id: 123 }, text: "held adoption" } }; + await writeTelegramSpooledUpdate({ + spoolDir: requireWebhookSpoolDir(), + update, + }); + let participant: TelegramSpooledReplayDeferredParticipant | undefined; + let settlementHold: TelegramSpooledReplaySettlementHold | undefined; + handleUpdateSpy.mockImplementationOnce(async () => { + participant = + createTelegramSpooledReplayDeferredParticipant("test:webhook-adoption-hold") ?? undefined; + settlementHold = participant?.beginSettlementHold(); + }); + + 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(participant).toBeDefined()); + await vi.advanceTimersByTimeAsync(25 * 60_000 + 10_000); + await yieldWebhookTask(); + + expect(participant?.abortSignal.aborted).toBe(false); + expect( + (await listTelegramSpooledUpdateClaims({ spoolDir: requireWebhookSpoolDir() })).map( + (claim) => claim.updateId, + ), + ).toEqual([42]); + + settlementHold?.release("discard-pending"); + participant?.settle({ kind: "completed" }); + await vi.waitFor(async () => + expect( + await listTelegramSpooledUpdateClaims({ spoolDir: requireWebhookSpoolDir() }), + ).toEqual([]), + ); + } finally { + settlementHold?.release("replay-pending"); + participant?.settle({ kind: "skipped" }); + 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({ diff --git a/extensions/telegram/src/webhook.ts b/extensions/telegram/src/webhook.ts index 823115ad334e..77ae84613de0 100644 --- a/extensions/telegram/src/webhook.ts +++ b/extensions/telegram/src/webhook.ts @@ -528,31 +528,34 @@ async function waitForWebhookSpooledDeferredWork(params: { log: (line: string) => void; update: ClaimedTelegramSpooledUpdate; }): Promise { - let timer: ReturnType | undefined; - const timeout = new Promise((resolve) => { - timer = setTimeout(() => { - const age = formatDurationPrecise(TELEGRAM_WEBHOOK_SPOOLED_HANDLER_TIMEOUT_MS); - const message = `Telegram webhook spool buffered processing timed out behind update ${params.update.updateId} on lane ${params.laneKey} after ${age}; marking the update failed.`; - params.log(`[telegram] ${message}`); - params.deferredWork.settle({ - kind: "failed-retryable", - error: new Error(message), - }); - resolve({ kind: "failed-retryable", error: new Error(message), timedOut: true }); - }, TELEGRAM_WEBHOOK_SPOOLED_HANDLER_TIMEOUT_MS); - timer.unref?.(); - }); + let timeoutError: Error | undefined; + const timer = setTimeout(() => { + const age = formatDurationPrecise(TELEGRAM_WEBHOOK_SPOOLED_HANDLER_TIMEOUT_MS); + const message = `Telegram webhook spool buffered processing timed out behind update ${params.update.updateId} on lane ${params.laneKey} after ${age}; marking the update failed.`; + params.log(`[telegram] ${message}`); + timeoutError = new Error(message); + params.deferredWork.settle({ + kind: "failed-retryable", + error: timeoutError, + }); + }, TELEGRAM_WEBHOOK_SPOOLED_HANDLER_TIMEOUT_MS); + timer.unref?.(); try { - return await Promise.race([ - params.deferredWork.task.catch((err: unknown): TelegramMessageProcessingResult => { - return { kind: "failed-retryable", error: err }; + const result = await params.deferredWork.task.catch( + (err: unknown): TelegramMessageProcessingResult => ({ + kind: "failed-retryable", + error: err, }), - timeout, - ]); + ); + // A durable-adoption hold can discard the timeout and settle completed. + // Only the exact timeout result owns the handler-timeout failure path. + return timeoutError !== undefined && + result.kind === "failed-retryable" && + result.error === timeoutError + ? { ...result, timedOut: true } + : result; } finally { - if (timer) { - clearTimeout(timer); - } + clearTimeout(timer); } }