fix(telegram): serialize spool timeout with adoption

(cherry picked from commit bdbac2a7e3)
This commit is contained in:
Vincent Koc
2026-07-10 05:22:26 -07:00
committed by Vincent Koc
parent 76bfe206a5
commit 4370233f0c
9 changed files with 252 additions and 58 deletions
+38 -16
View File
@@ -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<TelegramSpooledReplaySettlementHold["release"]>[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<TelegramMessageProcessingResult> | 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);
+1 -4
View File
@@ -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());
@@ -18,9 +18,15 @@ export type TelegramSpooledReplayDeferredParticipant = {
key: string;
abortSignal: AbortSignal;
task: Promise<TelegramMessageProcessingResult>;
/** 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<TelegramUpdateProcessingFrame>();
const telegramSpooledReplayFrames = new AsyncLocalStorage<TelegramSpooledReplayFrame>();
const telegramSpooledReplayUpdates = new WeakSet<object>();
@@ -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<TelegramMessageProcessingResult>((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);
},
};
}
@@ -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 {
@@ -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);
}
}
}
@@ -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();
@@ -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),
+58
View File
@@ -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({
+25 -22
View File
@@ -528,31 +528,34 @@ async function waitForWebhookSpooledDeferredWork(params: {
log: (line: string) => void;
update: ClaimedTelegramSpooledUpdate;
}): Promise<WebhookSpooledDeferredWorkResult> {
let timer: ReturnType<typeof setTimeout> | undefined;
const timeout = new Promise<WebhookSpooledDeferredWorkResult>((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);
}
}