From af24487b5f962a203fdab74ebd3bc876a7b89ec9 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Fri, 17 Jul 2026 23:10:54 +0100 Subject: [PATCH] refactor(signal): remove global ingress test hook (#110094) * test(signal): remove global ingress test hook * fix(signal): remove obsolete ingress test helpers --- extensions/signal/src/runtime.test-support.ts | 43 --- extensions/signal/src/signal-ingress.test.ts | 258 ++++++++++-------- extensions/signal/src/signal-ingress.ts | 33 +-- 3 files changed, 157 insertions(+), 177 deletions(-) diff --git a/extensions/signal/src/runtime.test-support.ts b/extensions/signal/src/runtime.test-support.ts index 9d06326fae21..8309d3e68b77 100644 --- a/extensions/signal/src/runtime.test-support.ts +++ b/extensions/signal/src/runtime.test-support.ts @@ -1,49 +1,6 @@ // Signal test support owns cleanup for process-global plugin runtime state. -import type { - ChannelIngressDrain, - ChannelIngressQueue, -} from "openclaw/plugin-sdk/channel-outbound"; import type { PluginRuntime } from "openclaw/plugin-sdk/core"; import { createPluginRuntimeStore } from "openclaw/plugin-sdk/runtime-store"; -import type { SignalSseEvent } from "./client-adapter.js"; -import "./signal-ingress.js"; - -export type SignalIngressPayload = { - version: 1; - receivedAt: number; - event: SignalSseEvent; -}; - -type SignalIngressEnqueueResult = - | Awaited["enqueue"]>> - | { kind: "ignored" }; - -type SignalIngressTestApi = { - createSignalIngressDrain(params: { - queue: ChannelIngressQueue; - dispatch: ( - event: SignalSseEvent, - lifecycle: { - abortSignal: AbortSignal; - onAdopted: () => void | Promise; - onDeferred: () => void; - onAdoptionFinalizing: () => void; - onAbandoned: () => void; - }, - ) => unknown; - }): ChannelIngressDrain; - enqueueSignalIngressEvent(params: { - queue: ChannelIngressQueue; - event: SignalSseEvent; - now?: number; - }): Promise; - resolveSignalIngressEventId(event: SignalSseEvent): string | null; - resolveSignalIngressLaneKey(event: SignalSseEvent): string | null; -}; - -export const signalIngressTesting = (globalThis as Record)[ - Symbol.for("openclaw.signalIngressTestApi") -] as SignalIngressTestApi; const { clearRuntime } = createPluginRuntimeStore({ pluginId: "signal", diff --git a/extensions/signal/src/signal-ingress.test.ts b/extensions/signal/src/signal-ingress.test.ts index 37c3de85212e..add5686ffb58 100644 --- a/extensions/signal/src/signal-ingress.test.ts +++ b/extensions/signal/src/signal-ingress.test.ts @@ -2,34 +2,55 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; -import type { ChannelIngressQueue } from "openclaw/plugin-sdk/channel-outbound"; -import type { PluginRuntime } from "openclaw/plugin-sdk/core"; import { closeOpenClawStateDatabaseForTest, createChannelIngressQueueForTests, } from "openclaw/plugin-sdk/plugin-state-test-runtime"; import { afterEach, describe, expect, it, vi } from "vitest"; import type { SignalSseEvent } from "./client-adapter.js"; -import { setSignalRuntime } from "./runtime.js"; -import { - clearSignalRuntimeForTest, - signalIngressTesting, - type SignalIngressPayload, -} from "./runtime.test-support.js"; import { startSignalIngressMonitor } from "./signal-ingress.js"; -const createSignalIngressDrain = ( - ...args: Parameters -) => signalIngressTesting.createSignalIngressDrain(...args); -const enqueueSignalIngressEvent = ( - ...args: Parameters -) => signalIngressTesting.enqueueSignalIngressEvent(...args); -const resolveSignalIngressEventId = ( - ...args: Parameters -) => signalIngressTesting.resolveSignalIngressEventId(...args); -const resolveSignalIngressLaneKey = ( - ...args: Parameters -) => signalIngressTesting.resolveSignalIngressLaneKey(...args); +type SignalIngressQueue = NonNullable[0]["queue"]>; +type SignalIngressPayload = Parameters[1]; +type SignalIngressDispatch = Parameters[0]["dispatch"]; + +function createTrackedTaskRunner() { + const pending = new Set>(); + const failures: unknown[] = []; + + const runTrackedTask = (task: () => Promise) => { + const promise = task() + .catch((error: unknown) => { + failures.push(error); + }) + .finally(() => { + pending.delete(promise); + }); + pending.add(promise); + }; + const waitForIdle = async () => { + while (pending.size > 0) { + await Promise.all(pending); + } + if (failures.length > 0) { + throw failures[0]; + } + }; + + return { runTrackedTask, waitForIdle }; +} + +async function startMonitor(queue: SignalIngressQueue, dispatch: SignalIngressDispatch) { + const tasks = createTrackedTaskRunner(); + const monitor = await startSignalIngressMonitor({ + accountId: "default", + queue, + dispatch, + runtime: { error: vi.fn(), log: vi.fn() }, + runTrackedTask: tasks.runTrackedTask, + }); + return { monitor, waitForIdle: tasks.waitForIdle }; +} function signalEvent(params?: { senderNumber?: string; @@ -57,7 +78,7 @@ function signalEvent(params?: { } async function withQueue( - fn: (queue: ChannelIngressQueue, stateDir: string) => Promise, + fn: (queue: SignalIngressQueue, stateDir: string) => Promise, ): Promise { const created = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-signal-ingress-")); const stateDir = await fs.realpath(created); @@ -75,7 +96,6 @@ async function withQueue( } afterEach(() => { - clearSignalRuntimeForTest(); closeOpenClawStateDatabaseForTest(); vi.restoreAllMocks(); }); @@ -87,13 +107,11 @@ describe("Signal durable ingress", () => { const failingQueue = { ...queue, enqueue: vi.fn().mockRejectedValue(appendError), - } satisfies ChannelIngressQueue; - setSignalRuntime({ - state: { openChannelIngressQueue: () => failingQueue }, - } as unknown as PluginRuntime); + } satisfies SignalIngressQueue; const dispatch = vi.fn(); const monitor = await startSignalIngressMonitor({ accountId: "default", + queue: failingQueue, dispatch, runtime: { error: vi.fn(), log: vi.fn() }, runTrackedTask: vi.fn(), @@ -110,69 +128,72 @@ describe("Signal durable ingress", () => { it("recovers an uncompleted append with a fresh drain and dispatches exactly once", async () => { await withQueue(async (queue) => { const event = signalEvent(); - await enqueueSignalIngressEvent({ queue, event }); + const interruptedDispatch = vi.fn((_event, lifecycle) => { + lifecycle.onDeferred(); + return { kind: "deferred" } as const; + }); + const interrupted = await startMonitor(queue, interruptedDispatch); + await interrupted.monitor.receive(event); + await interrupted.waitForIdle(); + expect(await queue.listClaims()).toHaveLength(1); + await interrupted.monitor.stop(); - const dispatch = vi.fn().mockResolvedValue(undefined); - const recoveredDrain = createSignalIngressDrain({ queue, dispatch }); - await recoveredDrain.drainOnce(); - await recoveredDrain.waitForIdle(); - recoveredDrain.dispose(); - - const restartedDrain = createSignalIngressDrain({ queue, dispatch }); - await restartedDrain.drainOnce(); - await restartedDrain.waitForIdle(); - restartedDrain.dispose(); - - expect(dispatch).toHaveBeenCalledTimes(1); - expect(dispatch).toHaveBeenCalledWith(event, expect.any(Object)); + const recoveredDispatch = vi.fn().mockResolvedValue(undefined); + const recovered = await startMonitor(queue, recoveredDispatch); + try { + await recovered.waitForIdle(); + expect(recoveredDispatch).toHaveBeenCalledTimes(1); + expect(recoveredDispatch).toHaveBeenCalledWith(event, expect.any(Object)); + } finally { + await recovered.monitor.stop(); + } }); }); it("keeps a completion tombstone so a duplicate cannot dispatch twice", async () => { await withQueue(async (queue) => { const event = signalEvent(); - const first = await enqueueSignalIngressEvent({ queue, event }); - expect(first.kind).toBe("accepted"); - const dispatch = vi.fn().mockResolvedValue(undefined); - const drain = createSignalIngressDrain({ queue, dispatch }); - await drain.drainOnce(); - await drain.waitForIdle(); - - const duplicate = await enqueueSignalIngressEvent({ queue, event }); - expect(duplicate.kind).toBe("completed"); - await drain.drainOnce(); - await drain.waitForIdle(); - drain.dispose(); - - expect(dispatch).toHaveBeenCalledTimes(1); + const started = await startMonitor(queue, dispatch); + try { + await started.monitor.receive(event); + await started.waitForIdle(); + await started.monitor.receive(event); + await started.waitForIdle(); + expect(dispatch).toHaveBeenCalledTimes(1); + } finally { + await started.monitor.stop(); + } }); }); it("completes only when deferred dispatch adoption becomes durable", async () => { await withQueue(async (queue) => { const event = signalEvent(); - await enqueueSignalIngressEvent({ queue, event }); let adopt: (() => void | Promise) | undefined; - const drain = createSignalIngressDrain({ - queue, - dispatch: (_event, lifecycle) => { - adopt = lifecycle.onAdopted; - lifecycle.onDeferred(); - return { kind: "deferred" }; - }, + const dispatch = vi.fn((_event, lifecycle) => { + adopt = lifecycle.onAdopted; + lifecycle.onDeferred(); + return { kind: "deferred" } as const; }); - - await drain.drainOnce(); - await vi.waitFor(async () => { + const started = await startMonitor(queue, dispatch); + try { + await started.monitor.receive(event); + await started.waitForIdle(); expect(await queue.listClaims()).toHaveLength(1); - }); - expect((await enqueueSignalIngressEvent({ queue, event })).kind).toBe("claimed"); - await adopt?.(); - await drain.waitForIdle(); - expect((await enqueueSignalIngressEvent({ queue, event })).kind).toBe("completed"); - drain.dispose(); + await started.monitor.receive(event); + await started.waitForIdle(); + expect(dispatch).toHaveBeenCalledTimes(1); + + expect(adopt).toBeDefined(); + await adopt?.(); + await started.monitor.receive(event); + await started.waitForIdle(); + expect(dispatch).toHaveBeenCalledTimes(1); + } finally { + await started.monitor.stop(); + } }); }); @@ -188,16 +209,16 @@ describe("Signal durable ingress", () => { { receivedAt: 1, laneKey: "direct:number:+15550001111" }, ); const dispatch = vi.fn(); - const drain = createSignalIngressDrain({ queue, dispatch }); - - await drain.drainOnce(); - await drain.waitForIdle(); - - expect((await queue.enqueue("malformed-event", {} as SignalIngressPayload)).kind).toBe( - "failed", - ); - expect(dispatch).not.toHaveBeenCalled(); - drain.dispose(); + const started = await startMonitor(queue, dispatch); + try { + await started.waitForIdle(); + expect((await queue.enqueue("malformed-event", {} as SignalIngressPayload)).kind).toBe( + "failed", + ); + expect(dispatch).not.toHaveBeenCalled(); + } finally { + await started.monitor.stop(); + } }); }); @@ -215,21 +236,17 @@ describe("Signal durable ingress", () => { timestamp: 1_700_000_000_099, message: "redelivered message", }); - expect(resolveSignalIngressEventId(original)).toBe(resolveSignalIngressEventId(redelivery)); - - await enqueueSignalIngressEvent({ queue, event: original }); const dispatch = vi.fn().mockResolvedValue(undefined); - const drain = createSignalIngressDrain({ queue, dispatch }); - await drain.drainOnce(); - await drain.waitForIdle(); - - const duplicate = await enqueueSignalIngressEvent({ queue, event: redelivery }); - expect(duplicate.kind).toBe("completed"); - await drain.drainOnce(); - await drain.waitForIdle(); - drain.dispose(); - - expect(dispatch).toHaveBeenCalledTimes(1); + const started = await startMonitor(queue, dispatch); + try { + await started.monitor.receive(original); + await started.waitForIdle(); + await started.monitor.receive(redelivery); + await started.waitForIdle(); + expect(dispatch).toHaveBeenCalledTimes(1); + } finally { + await started.monitor.stop(); + } }); }); @@ -237,17 +254,31 @@ describe("Signal durable ingress", () => { await withQueue(async (queue) => { const direct = signalEvent({ senderUuid: "123e4567-e89b-12d3-a456-426614174000" }); const group = signalEvent({ groupId: "group-123" }); + const dispatch = vi.fn((_event, lifecycle) => { + lifecycle.onDeferred(); + return { kind: "deferred" } as const; + }); + const started = await startMonitor(queue, dispatch); + try { + await started.monitor.receive(direct); + await started.monitor.receive(group); + await started.waitForIdle(); - expect(resolveSignalIngressLaneKey(direct)).toBe( - "direct:uuid:123e4567-e89b-12d3-a456-426614174000", - ); - expect(resolveSignalIngressLaneKey(group)).toBe("group:group-123"); - - await enqueueSignalIngressEvent({ queue, event: group }); - const pending = await queue.listPending({ limit: "all" }); - expect(pending).toHaveLength(1); - expect(pending[0]?.payload.event).toEqual(group); - expect(pending[0]?.laneKey).toBe("group:group-123"); + expect(await queue.listClaims()).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + laneKey: "direct:uuid:123e4567-e89b-12d3-a456-426614174000", + payload: expect.objectContaining({ event: direct }), + }), + expect.objectContaining({ + laneKey: "group:group-123", + payload: expect.objectContaining({ event: group }), + }), + ]), + ); + } finally { + await started.monitor.stop(); + } }); }); @@ -257,12 +288,17 @@ describe("Signal durable ingress", () => { ["typing", { envelope: { sourceNumber: "+15550001111", timestamp: 3, typingMessage: {} } }], ])("does not journal %s envelopes", async (_label, payload) => { await withQueue(async (queue) => { - const result = await enqueueSignalIngressEvent({ - queue, - event: { event: "receive", data: JSON.stringify(payload) }, - }); - expect(result.kind).toBe("ignored"); - await expect(queue.listPending({ limit: "all" })).resolves.toHaveLength(0); + const dispatch = vi.fn(); + const started = await startMonitor(queue, dispatch); + try { + await started.monitor.receive({ event: "receive", data: JSON.stringify(payload) }); + await started.waitForIdle(); + await expect(queue.listPending({ limit: "all" })).resolves.toHaveLength(0); + await expect(queue.listClaims()).resolves.toHaveLength(0); + expect(dispatch).not.toHaveBeenCalled(); + } finally { + await started.monitor.stop(); + } }); }); }); diff --git a/extensions/signal/src/signal-ingress.ts b/extensions/signal/src/signal-ingress.ts index 984fde3408a1..6b60100517cb 100644 --- a/extensions/signal/src/signal-ingress.ts +++ b/extensions/signal/src/signal-ingress.ts @@ -152,14 +152,6 @@ function inspectSignalIngressEvent(event: SignalSseEvent): SignalIngressEventFac }; } -function resolveSignalIngressEventId(event: SignalSseEvent): string | null { - return inspectSignalIngressEvent(event)?.eventId ?? null; -} - -function resolveSignalIngressLaneKey(event: SignalSseEvent): string | null { - return inspectSignalIngressEvent(event)?.laneKey ?? null; -} - type SignalIngressEnqueueResult = | Awaited["enqueue"]>> | { kind: "ignored" }; @@ -230,17 +222,21 @@ export type SignalIngressMonitor = { /** Open the account queue, recover it, and keep newly appended rows draining. */ export async function startSignalIngressMonitor(params: { accountId: string; + queue?: ChannelIngressQueue; dispatch: SignalIngressDispatch; runtime: Pick; runTrackedTask: (task: () => Promise) => void; }): Promise { - const pluginRuntime = getOptionalSignalRuntime(); - if (!pluginRuntime) { - throw new Error("Signal runtime not initialized for durable ingress"); + let queue = params.queue; + if (!queue) { + const pluginRuntime = getOptionalSignalRuntime(); + if (!pluginRuntime) { + throw new Error("Signal runtime not initialized for durable ingress"); + } + queue = pluginRuntime.state.openChannelIngressQueue({ + accountId: params.accountId, + }); } - const queue = pluginRuntime.state.openChannelIngressQueue({ - accountId: params.accountId, - }); const drain = createSignalIngressDrain({ queue, dispatch: params.dispatch, @@ -297,12 +293,3 @@ export async function startSignalIngressMonitor(params: { }, }; } - -if (process.env.VITEST || process.env.NODE_ENV === "test") { - (globalThis as Record)[Symbol.for("openclaw.signalIngressTestApi")] = { - createSignalIngressDrain, - enqueueSignalIngressEvent, - resolveSignalIngressEventId, - resolveSignalIngressLaneKey, - }; -}