diff --git a/extensions/telegram/src/test-support/webhook-gateway.ts b/extensions/telegram/src/test-support/webhook-gateway.ts index cf3dac9a85f1..f2b0a7388e35 100644 --- a/extensions/telegram/src/test-support/webhook-gateway.ts +++ b/extensions/telegram/src/test-support/webhook-gateway.ts @@ -5,6 +5,8 @@ import { setActivePluginRegistry, } from "openclaw/plugin-sdk/plugin-test-runtime"; import { canonicalizeWebhookRouteKey } from "openclaw/plugin-sdk/webhook-ingress"; +import { vi } from "vitest"; +import * as telegramIngressFactory from "../telegram-ingress-drain-factory.js"; type StartWebhook = typeof import("../webhook.js").startTelegramWebhook; type StartWebhookOptions = Omit[0], "token" | "abortSignal">; @@ -80,20 +82,36 @@ export function createTelegramWebhookTestGateway(options: { startWebhook, withWebhook: async ( params: StartWebhookOptions, - run: (ctx: { server: Server; port: number }) => Promise, + run: (ctx: { + server: Server; + port: number; + ingress: ReturnType; + }) => Promise, ): Promise => { - const abort = new AbortController(); - const started = await startWebhook({ - token: options.token, - abortSignal: abort.signal, - ...options.queueScope(), - ...params, - }); + const createIngress = telegramIngressFactory.createTelegramTransportIngressMonitor; + let ingress: ReturnType | undefined; + const ingressFactory = vi + .spyOn(telegramIngressFactory, "createTelegramTransportIngressMonitor") + .mockImplementation((ingressParams) => (ingress = createIngress(ingressParams))); try { - return await run({ server, port: getServerPort(server) }); + const abort = new AbortController(); + const started = await startWebhook({ + token: options.token, + abortSignal: abort.signal, + ...options.queueScope(), + ...params, + }); + try { + if (!ingress) { + throw new Error("Expected the started webhook's ingress monitor"); + } + return await run({ server, port: getServerPort(server), ingress }); + } finally { + await started.stop(); + abort.abort(); + } } finally { - await started.stop(); - abort.abort(); + ingressFactory.mockRestore(); } }, }; diff --git a/extensions/telegram/src/webhook.test.ts b/extensions/telegram/src/webhook.test.ts index 460b5fffcc6d..d3c5a83c362a 100644 --- a/extensions/telegram/src/webhook.test.ts +++ b/extensions/telegram/src/webhook.test.ts @@ -38,7 +38,6 @@ import { installTelegramIngressQueueRuntime } from "./runtime-state.test-support import { setTelegramRuntime } from "./runtime.js"; import { clearTelegramRuntimeForTest as clearTelegramRuntime } from "./runtime.test-support.js"; import type { TelegramRuntime } from "./runtime.types.js"; -import * as telegramIngressFactory from "./telegram-ingress-drain-factory.js"; import { openTelegramIngressQueue } from "./telegram-ingress-spool.js"; import { writeTelegramSpooledUpdate, @@ -161,6 +160,7 @@ const { startWebhook: startTelegramWebhook, server: gatewayServer, pendingRequests: pendingRouteRequests, + withWebhook: withStartedWebhook, } = gateway; let webhookStateDir: string | undefined; @@ -199,6 +199,8 @@ beforeEach(async () => { resetTelegramWebhookMocks(); webhookStateDir = await fs.mkdtemp(nodePath.join(os.tmpdir(), "openclaw-telegram-webhook-")); installTelegramIngressQueueRuntime(() => webhookStateDir ?? os.tmpdir()); + // The production monitor prepares shared state before starting the webhook. + await openTelegramIngressQueue(requireWebhookQueueScope()).listPending(); }); afterEach(async () => { @@ -214,31 +216,6 @@ afterEach(async () => { } }); -async function withStartedWebhook( - options: Parameters[0], - run: (ctx: { - server: typeof gateway.server; - port: number; - ingress: ReturnType; - }) => Promise, -): Promise { - const createIngress = telegramIngressFactory.createTelegramTransportIngressMonitor; - let ingress: ReturnType | undefined; - const ingressFactory = vi - .spyOn(telegramIngressFactory, "createTelegramTransportIngressMonitor") - .mockImplementation((params) => (ingress = createIngress(params))); - try { - return await gateway.withWebhook(options, async (ctx) => { - if (!ingress) { - throw new Error("Expected the started webhook's ingress monitor"); - } - return await run({ ...ctx, ingress }); - }); - } finally { - ingressFactory.mockRestore(); - } -} - function startWebhookStartupFixture( options: Partial[0]> = {}, ) {