



























@@ -639,6 +639,105 @@ describe("TelegramPollingSession", () => {
639639}
640640});
641641642+it("lets isolated ingress drain interleave different Telegram chats", async () => {
643+const abort = new AbortController();
644+const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
645+const events: string[] = [];
646+let releaseFirstChatTurn: (() => void) | undefined;
647+const firstChatTurnDone = new Promise<void>((resolve) => {
648+releaseFirstChatTurn = resolve;
649+});
650+const handleUpdate = vi.fn(async (update: { update_id?: number }) => {
651+if (update.update_id === 42) {
652+events.push("chatA:start");
653+await firstChatTurnDone;
654+events.push("chatA:end");
655+return;
656+}
657+if (update.update_id === 43) {
658+events.push("chatB");
659+return;
660+}
661+if (update.update_id === 44) {
662+events.push("chatA:second");
663+}
664+});
665+const bot = {
666+api: {
667+deleteWebhook: vi.fn(async () => true),
668+config: { use: vi.fn() },
669+},
670+init: vi.fn(async () => undefined),
671+ handleUpdate,
672+stop: vi.fn(async () => undefined),
673+};
674+createTelegramBotMock.mockReturnValueOnce(bot);
675+for (const { updateId, chatId, text } of [
676+{ updateId: 42, chatId: -100, text: "long first chat turn" },
677+{ updateId: 43, chatId: 854067528, text: "second chat turn" },
678+{ updateId: 44, chatId: -100, text: "second first chat turn" },
679+]) {
680+await writeTelegramSpooledUpdate({
681+spoolDir: tempDir,
682+update: {
683+update_id: updateId,
684+message: {
685+ text,
686+chat: { id: chatId, type: chatId < 0 ? "supergroup" : "private" },
687+},
688+},
689+});
690+}
691+let stopWorker: (() => void) | undefined;
692+const workerDone = new Promise<void>((resolve) => {
693+stopWorker = resolve;
694+});
695+const createWorker = vi.fn(() => ({
696+onMessage: vi.fn(() => () => undefined),
697+stop: vi.fn(async () => {
698+stopWorker?.();
699+}),
700+task: vi.fn(async () => {
701+await workerDone;
702+}),
703+}));
704+705+try {
706+const session = createPollingSession({
707+abortSignal: abort.signal,
708+isolatedIngress: {
709+enabled: true,
710+spoolDir: tempDir,
711+ createWorker,
712+drainIntervalMs: 10,
713+},
714+});
715+716+const runPromise = session.runUntilAbort();
717+await vi.waitFor(() => expect(events).toEqual(["chatA:start", "chatB"]));
718+expect(
719+(await listTelegramSpooledUpdates({ spoolDir: tempDir })).map((update) => update.updateId),
720+).toEqual([42, 44]);
721+722+releaseFirstChatTurn?.();
723+await vi.waitFor(() =>
724+expect(events).toEqual(["chatA:start", "chatB", "chatA:end", "chatA:second"]),
725+);
726+await vi.waitFor(async () =>
727+expect(
728+(await listTelegramSpooledUpdates({ spoolDir: tempDir })).map(
729+(update) => update.updateId,
730+),
731+).toEqual([]),
732+);
733+abort.abort();
734+await runPromise;
735+} finally {
736+releaseFirstChatTurn?.();
737+await fs.rm(tempDir, { recursive: true, force: true });
738+}
739+});
740+642741it("lets isolated ingress control updates bypass an active spooled turn", async () => {
643742const abort = new AbortController();
644743const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。