Skip to content

Commit db255b1

Browse files
authored
Fix Telegram spooled claim refresh (#96962)
1 parent 4fc504d commit db255b1

6 files changed

Lines changed: 547 additions & 20 deletions

File tree

extensions/telegram/src/polling-session.test.ts

Lines changed: 230 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -486,6 +486,49 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
486486
return (await listTelegramSpooledUpdates({ spoolDir, limit })).map((update) => update.updateId);
487487
}
488488

489+
async function claimedAtForUpdate(spoolDir: string, updateId: number): Promise<number> {
490+
const claim = (await listTelegramSpooledUpdateClaims({ spoolDir })).find(
491+
(entry) => entry.updateId === updateId,
492+
);
493+
if (!claim?.claim) {
494+
throw new Error(`Expected claimed spooled update ${updateId}`);
495+
}
496+
return claim.claim.claimedAt;
497+
}
498+
499+
function installSpooledClaimRefreshHarness(): {
500+
restore: () => void;
501+
triggerRefresh: () => void;
502+
} {
503+
let refresh: (() => void) | undefined;
504+
const realSetInterval = globalThis.setInterval.bind(globalThis);
505+
const setIntervalSpy = vi.spyOn(globalThis, "setInterval").mockImplementation(((
506+
handler: Parameters<typeof setInterval>[0],
507+
timeout?: number,
508+
) => {
509+
if (timeout === pollingSessionTesting.spooledClaimRefreshIntervalMs) {
510+
refresh = () => {
511+
if (typeof handler === "function") {
512+
handler();
513+
}
514+
};
515+
const timer = realSetInterval(() => undefined, 2_147_483_647);
516+
timer.unref?.();
517+
return timer;
518+
}
519+
return realSetInterval(handler, timeout);
520+
}) as typeof setInterval);
521+
return {
522+
restore: () => setIntervalSpy.mockRestore(),
523+
triggerRefresh: () => {
524+
if (!refresh) {
525+
throw new Error("Expected spooled claim refresh interval to be registered");
526+
}
527+
refresh();
528+
},
529+
};
530+
}
531+
489532
function normalizeTelegramTestAccountId(spoolDir: string): string {
490533
const trimmed = path.basename(spoolDir).trim();
491534
return trimmed ? trimmed.replace(/[^a-z0-9._-]+/gi, "_") : "default";
@@ -1575,6 +1618,49 @@ describe("TelegramPollingSession", () => {
15751618
});
15761619
});
15771620

1621+
it("refreshes active spooled claims while the handler is still running", async () => {
1622+
const refreshHarness = installSpooledClaimRefreshHarness();
1623+
await withTempSpool(async (tempDir) => {
1624+
const abort = new AbortController();
1625+
const events: string[] = [];
1626+
let releaseHandler: (() => void) | undefined;
1627+
const handlerDone = new Promise<void>((resolve) => {
1628+
releaseHandler = resolve;
1629+
});
1630+
await writeSpooledTestUpdates(tempDir, [topicUpdate(42, 10, "long topic 10 turn")]);
1631+
1632+
const { runPromise, stopWorker } = startIsolatedIngressSession({
1633+
abort,
1634+
spoolDir: tempDir,
1635+
handleUpdate: async (update) => {
1636+
events.push(`topic10:${update.update_id}`);
1637+
await handlerDone;
1638+
},
1639+
});
1640+
1641+
try {
1642+
await vi.waitFor(() => expect(events).toEqual(["topic10:42"]));
1643+
const before = await claimedAtForUpdate(tempDir, 42);
1644+
1645+
refreshHarness.triggerRefresh();
1646+
await vi.waitFor(async () =>
1647+
expect(await claimedAtForUpdate(tempDir, 42)).toBeGreaterThan(before),
1648+
);
1649+
1650+
releaseHandler?.();
1651+
await vi.waitFor(async () =>
1652+
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]),
1653+
);
1654+
} finally {
1655+
releaseHandler?.();
1656+
abort.abort();
1657+
stopWorker();
1658+
refreshHarness.restore();
1659+
await runPromise;
1660+
}
1661+
});
1662+
});
1663+
15781664
it("holds buffered spooled claims until deferred processing settles without blocking same-lane buffering", async () => {
15791665
await withTempSpool(async (tempDir) => {
15801666
const abort = new AbortController();
@@ -1625,6 +1711,50 @@ describe("TelegramPollingSession", () => {
16251711
});
16261712
});
16271713

1714+
it("refreshes deferred spooled claims after the active handler hands off", async () => {
1715+
const refreshHarness = installSpooledClaimRefreshHarness();
1716+
await withTempSpool(async (tempDir) => {
1717+
const abort = new AbortController();
1718+
const participants: TelegramSpooledReplayDeferredParticipant[] = [];
1719+
await writeSpooledTestUpdates(tempDir, [topicUpdate(42, 10, "buffered topic 10 turn")]);
1720+
1721+
const { runPromise, stopWorker } = startIsolatedIngressSession({
1722+
abort,
1723+
spoolDir: tempDir,
1724+
handleUpdate: async (update) => {
1725+
const participant = createTelegramSpooledReplayDeferredParticipant(
1726+
`test-buffer:${update.update_id}`,
1727+
);
1728+
if (!participant) {
1729+
throw new Error("expected spooled replay participant");
1730+
}
1731+
participants.push(participant);
1732+
},
1733+
});
1734+
1735+
try {
1736+
await vi.waitFor(() => expect(participants).toHaveLength(1));
1737+
const before = await claimedAtForUpdate(tempDir, 42);
1738+
1739+
refreshHarness.triggerRefresh();
1740+
await vi.waitFor(async () =>
1741+
expect(await claimedAtForUpdate(tempDir, 42)).toBeGreaterThan(before),
1742+
);
1743+
1744+
participants[0]?.settle({ kind: "completed" });
1745+
await vi.waitFor(async () =>
1746+
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]),
1747+
);
1748+
} finally {
1749+
participants[0]?.settle({ kind: "completed" });
1750+
abort.abort();
1751+
stopWorker();
1752+
refreshHarness.restore();
1753+
await runPromise;
1754+
}
1755+
});
1756+
});
1757+
16281758
it("releases buffered spooled claims for retry when deferred processing fails", async () => {
16291759
await withTempSpool(async (tempDir) => {
16301760
const abort = new AbortController();
@@ -3585,6 +3715,106 @@ describe("TelegramPollingSession", () => {
35853715
}
35863716
});
35873717

3718+
it("marks isolated ingress unhealthy when a spooled backlog stalls before handler timeout", async () => {
3719+
vi.useFakeTimers({ now: 1_000, shouldAdvanceTime: true });
3720+
const abort = new AbortController();
3721+
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
3722+
const setStatus = vi.fn();
3723+
let releaseRegularTurn: (() => void) | undefined;
3724+
const regularTurnDone = new Promise<void>((resolve) => {
3725+
releaseRegularTurn = resolve;
3726+
});
3727+
const handleUpdate = vi.fn(async () => {
3728+
await regularTurnDone;
3729+
});
3730+
createTelegramBotMock.mockReturnValueOnce({
3731+
api: {
3732+
deleteWebhook: vi.fn(async () => true),
3733+
config: { use: vi.fn() },
3734+
},
3735+
init: vi.fn(async () => undefined),
3736+
handleUpdate,
3737+
stop: vi.fn(async () => undefined),
3738+
});
3739+
await writeSpooledTestUpdates(tempDir, [
3740+
topicUpdate(42, 10, "active topic 10 turn"),
3741+
topicUpdate(43, 10, "later topic 10 turn"),
3742+
]);
3743+
3744+
const workerListeners: WorkerMessageListener[] = [];
3745+
let stopWorker: (() => void) | undefined;
3746+
const workerDone = new Promise<void>((resolve) => {
3747+
stopWorker = resolve;
3748+
});
3749+
const createWorker = vi.fn(() => ({
3750+
onMessage: vi.fn((listener: WorkerMessageListener) => {
3751+
workerListeners.push(listener);
3752+
return () => undefined;
3753+
}),
3754+
stop: vi.fn(async () => {
3755+
stopWorker?.();
3756+
}),
3757+
task: vi.fn(async () => {
3758+
await workerDone;
3759+
}),
3760+
}));
3761+
3762+
try {
3763+
const session = createPollingSession({
3764+
abortSignal: abort.signal,
3765+
setStatus,
3766+
isolatedIngress: {
3767+
enabled: true,
3768+
spoolDir: tempDir,
3769+
createWorker,
3770+
drainIntervalMs: pollingSessionTesting.isolatedIngressBacklogStallMs * 2,
3771+
spooledUpdateHandlerTimeoutMs: pollingSessionTesting.isolatedIngressBacklogStallMs * 2,
3772+
},
3773+
});
3774+
3775+
const runPromise = session.runUntilAbort();
3776+
await vi.waitFor(() => expect(handleUpdate).toHaveBeenCalledTimes(1));
3777+
workerListeners[0]?.({
3778+
type: "poll-success",
3779+
offset: null,
3780+
count: 0,
3781+
finishedAt: Date.now(),
3782+
});
3783+
expect(statusPatches(setStatus).some((patch) => patch.connected === true)).toBe(true);
3784+
3785+
vi.setSystemTime(1_000 + pollingSessionTesting.isolatedIngressBacklogStallMs + 1);
3786+
workerListeners[0]?.({ type: "spooled", updateId: 43, queued: 1 });
3787+
await vi.waitFor(() =>
3788+
expect(
3789+
statusPatches(setStatus).some(
3790+
(patch) =>
3791+
patch.connected === false &&
3792+
String(patch.lastError).includes("isolated polling spool backlog stalled"),
3793+
),
3794+
).toBe(true),
3795+
);
3796+
expect(await failedUpdateIds(tempDir)).toEqual([]);
3797+
expect(await pendingUpdateIds(tempDir, "all")).toEqual([43]);
3798+
expect(
3799+
(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).map(
3800+
(claim) => claim.updateId,
3801+
),
3802+
).toEqual([42]);
3803+
3804+
releaseRegularTurn?.();
3805+
abort.abort();
3806+
stopWorker?.();
3807+
await vi.advanceTimersByTimeAsync(20_000);
3808+
await runPromise;
3809+
} finally {
3810+
releaseRegularTurn?.();
3811+
abort.abort();
3812+
stopWorker?.();
3813+
vi.useRealTimers();
3814+
await fs.rm(tempDir, { recursive: true, force: true });
3815+
}
3816+
});
3817+
35883818
it("marks isolated ingress unhealthy when a spooled backlog handler times out", async () => {
35893819
vi.useFakeTimers({ shouldAdvanceTime: true });
35903820
const abort = new AbortController();

0 commit comments

Comments
 (0)