


























@@ -486,6 +486,49 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
486486return (await listTelegramSpooledUpdates({ spoolDir, limit })).map((update) => update.updateId);
487487}
488488489+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+489532function normalizeTelegramTestAccountId(spoolDir: string): string {
490533const trimmed = path.basename(spoolDir).trim();
491534return trimmed ? trimmed.replace(/[^a-z0-9._-]+/gi, "_") : "default";
@@ -1575,6 +1618,49 @@ describe("TelegramPollingSession", () => {
15751618});
15761619});
157716201621+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+15781664it("holds buffered spooled claims until deferred processing settles without blocking same-lane buffering", async () => {
15791665await withTempSpool(async (tempDir) => {
15801666const abort = new AbortController();
@@ -1625,6 +1711,50 @@ describe("TelegramPollingSession", () => {
16251711});
16261712});
162717131714+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+16281758it("releases buffered spooled claims for retry when deferred processing fails", async () => {
16291759await withTempSpool(async (tempDir) => {
16301760const abort = new AbortController();
@@ -3585,6 +3715,106 @@ describe("TelegramPollingSession", () => {
35853715}
35863716});
358737173718+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+35883818it("marks isolated ingress unhealthy when a spooled backlog handler times out", async () => {
35893819vi.useFakeTimers({ shouldAdvanceTime: true });
35903820const abort = new AbortController();
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。