mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-27 21:07:01 -06:00
fix(telegram): keep timed-out webhook lanes guarded (#98806)
This commit is contained in:
@@ -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<void>((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();
|
||||
}
|
||||
|
||||
@@ -84,7 +84,11 @@ const TELEGRAM_WEBHOOK_REGISTRATION_RETRY_POLICY: BackoffPolicy = {
|
||||
factor: 2,
|
||||
jitter: 0.2,
|
||||
};
|
||||
const activeWebhookSpooledHandlersByLane = new Set<string>();
|
||||
type ActiveWebhookSpooledHandler = {
|
||||
laneKey: string;
|
||||
};
|
||||
|
||||
const activeWebhookSpooledHandlersByLane = new Map<string, ActiveWebhookSpooledHandler>();
|
||||
|
||||
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<string> {
|
||||
const laneKeys = new Set<string>();
|
||||
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<void> {
|
||||
}): Promise<boolean> {
|
||||
let timer: ReturnType<typeof setTimeout> | 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<void>((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<boolean>((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<unknown>;
|
||||
};
|
||||
|
||||
async function runWebhookSpooledReplayWithTimeout(params: {
|
||||
bot: ReturnType<typeof createTelegramBot>;
|
||||
laneKey: string;
|
||||
@@ -549,7 +561,7 @@ async function handleWebhookSpooledUpdate(params: {
|
||||
bot: ReturnType<typeof createTelegramBot>;
|
||||
log: (line: string) => void;
|
||||
update: ClaimedTelegramSpooledUpdate;
|
||||
}): Promise<void> {
|
||||
}): Promise<WebhookSpooledUpdateHandlerResult> {
|
||||
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<unknown> | 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;
|
||||
|
||||
Reference in New Issue
Block a user