



























@@ -3,6 +3,7 @@ import os from "node:os";
33import path from "node:path";
44import type { ChannelAccountSnapshot } from "openclaw/plugin-sdk/channel-contract";
55import { beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
6+import type { TelegramIngressWorkerMessage } from "./telegram-ingress-worker.js";
6778const runMock = vi.hoisted(() => vi.fn());
89const createTelegramBotMock = vi.hoisted(() => vi.fn());
@@ -90,6 +91,7 @@ type WorkerPollErrorListener = (message: {
9091message: string;
9192finishedAt: number;
9293}) => void;
94+type WorkerMessageListener = (message: TelegramIngressWorkerMessage) => void;
9395type AsyncVoidFn = () => Promise<void>;
9496type MockCallSource = { mock: { calls: Array<Array<unknown>> } };
9597@@ -176,6 +178,7 @@ function installPollingStallWatchdogHarness(dateNowSequence: readonly number[] =
176178break;
177179}
178180await Promise.resolve();
181+await new Promise<void>((resolve) => setImmediate(resolve));
179182}
180183expect(watchdog).toBeTypeOf("function");
181184return watchdog;
@@ -776,6 +779,142 @@ describe("TelegramPollingSession", () => {
776779await runPromise;
777780});
778781782+it("restarts isolated ingress when worker liveness stalls", async () => {
783+const abort = new AbortController();
784+const log = vi.fn();
785+const bot = {
786+api: {
787+deleteWebhook: vi.fn(async () => true),
788+config: { use: vi.fn() },
789+},
790+init: vi.fn(async () => undefined),
791+handleUpdate: vi.fn(async () => undefined),
792+stop: vi.fn(async () => undefined),
793+};
794+createTelegramBotMock.mockReturnValue(bot);
795+796+let firstWorkerDone: (() => void) | undefined;
797+const firstWorkerTask = new Promise<void>((resolve) => {
798+firstWorkerDone = resolve;
799+});
800+const firstWorkerStop = vi.fn(async () => {
801+firstWorkerDone?.();
802+});
803+let workerCycle = 0;
804+const createWorker = vi.fn(() => {
805+workerCycle += 1;
806+if (workerCycle === 1) {
807+return {
808+onMessage: vi.fn(() => () => undefined),
809+stop: firstWorkerStop,
810+task: vi.fn(async () => {
811+await firstWorkerTask;
812+}),
813+};
814+}
815+return {
816+onMessage: vi.fn(() => () => undefined),
817+stop: vi.fn(async () => undefined),
818+task: vi.fn(async () => {
819+abort.abort();
820+}),
821+};
822+});
823+const watchdogHarness = installPollingStallWatchdogHarness([0]);
824+const session = createPollingSession({
825+abortSignal: abort.signal,
826+ log,
827+stallThresholdMs: 30_000,
828+isolatedIngress: {
829+enabled: true,
830+ createWorker,
831+drainIntervalMs: 500,
832+},
833+});
834+835+try {
836+const runPromise = session.runUntilAbort();
837+const watchdog = await watchdogHarness.waitForWatchdog();
838+watchdogHarness.setNow(31_000);
839+watchdog?.();
840+841+await vi.waitFor(() => expect(firstWorkerStop).toHaveBeenCalledTimes(1));
842+await vi.waitFor(() => expect(createWorker).toHaveBeenCalledTimes(2));
843+await runPromise;
844+845+expectLogIncludes(log, "Polling stall detected");
846+expectLogIncludes(log, "isolated polling ingress finished reason=polling stall detected");
847+} finally {
848+watchdogHarness.restore();
849+abort.abort();
850+}
851+});
852+853+it("keeps isolated ingress alive when spooled messages show worker activity", async () => {
854+const abort = new AbortController();
855+const log = vi.fn();
856+const bot = {
857+api: {
858+deleteWebhook: vi.fn(async () => true),
859+config: { use: vi.fn() },
860+},
861+init: vi.fn(async () => undefined),
862+handleUpdate: vi.fn(async () => undefined),
863+stop: vi.fn(async () => undefined),
864+};
865+createTelegramBotMock.mockReturnValue(bot);
866+867+let onMessage: WorkerMessageListener | undefined;
868+let stopWorker: (() => void) | undefined;
869+const workerDone = new Promise<void>((resolve) => {
870+stopWorker = resolve;
871+});
872+const workerStop = vi.fn(async () => {
873+stopWorker?.();
874+});
875+const createWorker = vi.fn(() => ({
876+onMessage: vi.fn((handler: WorkerMessageListener) => {
877+onMessage = handler;
878+return () => undefined;
879+}),
880+stop: workerStop,
881+task: vi.fn(async () => {
882+await workerDone;
883+}),
884+}));
885+const watchdogHarness = installPollingStallWatchdogHarness([0]);
886+const session = createPollingSession({
887+abortSignal: abort.signal,
888+ log,
889+stallThresholdMs: 30_000,
890+isolatedIngress: {
891+enabled: true,
892+ createWorker,
893+drainIntervalMs: 500,
894+},
895+});
896+897+try {
898+const runPromise = session.runUntilAbort();
899+const watchdog = await watchdogHarness.waitForWatchdog();
900+onMessage?.({ type: "poll-start", offset: null, startedAt: 0 });
901+watchdogHarness.setNow(31_000);
902+onMessage?.({ type: "spooled", updateId: 42, queued: 1 });
903+watchdogHarness.setNow(45_000);
904+watchdog?.();
905+906+expect(workerStop).not.toHaveBeenCalled();
907+expectLogExcludes(log, "Polling stall detected");
908+909+abort.abort();
910+stopWorker?.();
911+await runPromise;
912+} finally {
913+watchdogHarness.restore();
914+abort.abort();
915+}
916+});
917+779918it("keeps failed lanes blocked for the rest of the drain pass", async () => {
780919await withTempSpool(async (tempDir) => {
781920const abort = new AbortController();
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。