diff --git a/docs/concepts/qa-e2e-automation.md b/docs/concepts/qa-e2e-automation.md index 51f26b389f15..92e612c92cde 100644 --- a/docs/concepts/qa-e2e-automation.md +++ b/docs/concepts/qa-e2e-automation.md @@ -1133,7 +1133,9 @@ The minimum adoption bar for a new channel: competing root command. Runner plugins should declare `qaRunners` in `openclaw.plugin.json` and export a matching `qaRunnerCliRegistrations` array from `runtime-api.ts`. Keep `runtime-api.ts` light; lazy CLI and - runner execution should stay behind separate entrypoints. + runner execution should stay behind separate entrypoints. An optional + `adapterFactory` exposes the transport to shared scenarios without changing + the command's existing scenario catalog. 5. Author or adapt YAML scenarios under the themed `qa/scenarios/` directories. 6. Use the generic scenario helpers for new scenarios. diff --git a/docs/help/testing.md b/docs/help/testing.md index cc24a30a0eb5..540de332c986 100644 --- a/docs/help/testing.md +++ b/docs/help/testing.md @@ -583,8 +583,9 @@ for Slack rows. The architecture and scenario-helper names for new channel adapters live in [QA overview - Adding a channel](/concepts/qa-e2e-automation#adding-a-channel). The minimum bar: implement the transport runner on the shared `qa-lab` host -seam, declare `qaRunners` in the plugin manifest, mount as -`openclaw qa `, and author scenarios under `qa/scenarios/`. +seam, add an `adapterFactory` for shared scenarios, declare `qaRunners` in the +plugin manifest, mount as `openclaw qa `, and author scenarios under +`qa/scenarios/`. ## Test suites (what runs where) diff --git a/docs/plugins/manifest.md b/docs/plugins/manifest.md index 43912d273a1d..48938c8115bb 100644 --- a/docs/plugins/manifest.md +++ b/docs/plugins/manifest.md @@ -405,7 +405,12 @@ Planner diagnostics can distinguish explicit activation hints from manifest owne ## qaRunners reference -Use `qaRunners` when a plugin contributes one or more transport runners beneath the shared `openclaw qa` root. Keep this metadata cheap and static; the plugin runtime still owns actual CLI registration through a lightweight `runtime-api.ts` surface that exports `qaRunnerCliRegistrations`. +Use `qaRunners` when a plugin contributes one or more transport runners beneath +the shared `openclaw qa` root. Keep this metadata cheap and static; the plugin +runtime still owns actual CLI registration through a lightweight +`runtime-api.ts` surface that exports matching `qaRunnerCliRegistrations`. An +optional `adapterFactory` exposes the transport to shared QA scenarios without +changing the registered command's runner. ```json { @@ -423,6 +428,9 @@ Use `qaRunners` when a plugin contributes one or more transport runners beneath | `commandName` | Yes | `string` | Subcommand mounted beneath `openclaw qa`, for example `matrix`. | | `description` | No | `string` | Fallback help text used when the shared host needs a stub command. | +The `adapterFactory` id must match `commandName`. Do not export registrations +for commands absent from the manifest. + ## setup reference Use `setup` when setup and onboarding surfaces need cheap plugin-owned metadata before runtime loads. diff --git a/extensions/qa-lab/src/cli.runtime.test.ts b/extensions/qa-lab/src/cli.runtime.test.ts index 1076dbef70a4..1a830aa75365 100644 --- a/extensions/qa-lab/src/cli.runtime.test.ts +++ b/extensions/qa-lab/src/cli.runtime.test.ts @@ -10,6 +10,7 @@ const { runQaSuite, runQaCharacterEval, runQaMultipass, + listLiveTransportQaAdapterFactories, listTelegramQaScenarioCatalog, runTelegramQaLive, startQaLabServer, @@ -23,6 +24,7 @@ const { runQaSuite: vi.fn(), runQaCharacterEval: vi.fn(), runQaMultipass: vi.fn(), + listLiveTransportQaAdapterFactories: vi.fn(), listTelegramQaScenarioCatalog: vi.fn(), runTelegramQaLive: vi.fn(), startQaLabServer: vi.fn(), @@ -50,6 +52,10 @@ vi.mock("./multipass.runtime.js", () => ({ runQaMultipass, })); +vi.mock("./live-transports/cli.js", () => ({ + listLiveTransportQaAdapterFactories, +})); + vi.mock("./live-transports/telegram/telegram-live.runtime.js", () => ({ listTelegramQaScenarioCatalog, runTelegramQaLive, @@ -222,6 +228,7 @@ describe("qa cli runtime", () => { runQaCharacterEval.mockReset(); runQaManualLane.mockReset(); runQaMultipass.mockReset(); + listLiveTransportQaAdapterFactories.mockReset(); listTelegramQaScenarioCatalog.mockReset(); runTelegramQaLive.mockReset(); startQaLabServer.mockReset(); @@ -282,6 +289,14 @@ describe("qa cli runtime", () => { regressionRefs: ["openclaw/openclaw#74698"], }, ]); + listLiveTransportQaAdapterFactories.mockReturnValue([ + { + id: "telegram", + scenarioIds: ["channel-chat-baseline"], + matches: vi.fn(), + create: vi.fn(), + }, + ]); startQaLabServer.mockResolvedValue({ baseUrl: "http://127.0.0.1:58000", runSelfCheck: vi.fn().mockResolvedValue({ @@ -644,6 +659,77 @@ describe("qa cli runtime", () => { }); }); + it("runs canonical scenarios through a discovered live adapter factory", async () => { + await runQaSuiteCommand({ + repoRoot: "/tmp/openclaw-repo", + outputDir: ".artifacts/qa/telegram-live", + channelDriver: "live", + channel: "telegram", + providerMode: "mock-openai", + scenarioIds: ["channel-chat-baseline"], + }); + + expect(runQaSuite).toHaveBeenCalledWith( + expect.objectContaining({ + adapterFactories: listLiveTransportQaAdapterFactories.mock.results[0]?.value, + channelDriver: "live", + channelId: "telegram", + concurrency: 1, + adapterOptions: expect.objectContaining({ + repoRoot: path.resolve("/tmp/openclaw-repo"), + }), + scenarioIds: ["channel-chat-baseline"], + }), + ); + }); + + it("uses the selected live adapter's declared scenarios by default", async () => { + await runQaSuiteCommand({ + channelDriver: "live", + channel: "telegram", + }); + + expect(runQaSuite).toHaveBeenCalledWith( + expect.objectContaining({ + scenarioIds: ["channel-chat-baseline"], + }), + ); + }); + + it("rejects live adapter selection under Multipass", async () => { + await expect( + runQaSuiteCommand({ + runner: "multipass", + channelDriver: "live", + channel: "telegram", + scenarioIds: ["channel-chat-baseline"], + }), + ).rejects.toThrow("--channel-driver live with --channel requires --runner host."); + expect(runQaMultipass).not.toHaveBeenCalled(); + }); + + it("rejects runtime-pair execution for live adapters", async () => { + await expect( + runQaSuiteCommand({ + channelDriver: "live", + channel: "telegram", + runtimePair: "openclaw,codex", + }), + ).rejects.toThrow("--runtime-pair is not supported with a live QA adapter."); + expect(runQaSuite).not.toHaveBeenCalled(); + }); + + it("keeps live taxonomy metadata unchanged without an explicit adapter channel", async () => { + await runQaSuiteCommand({ + channelDriver: "live", + scenarioIds: ["channel-chat-baseline"], + }); + + expect(runQaSuite).toHaveBeenCalledWith( + expect.not.objectContaining({ adapterFactories: expect.anything() }), + ); + }); + it("uses the Crabline default channel when selected scenarios do not request one", async () => { await runQaSuiteCommand({ repoRoot: "/tmp/openclaw-repo", diff --git a/extensions/qa-lab/src/cli.runtime.ts b/extensions/qa-lab/src/cli.runtime.ts index 585e187a711f..2e1b719edebf 100644 --- a/extensions/qa-lab/src/cli.runtime.ts +++ b/extensions/qa-lab/src/cli.runtime.ts @@ -43,8 +43,9 @@ import { type JsonlReplayInput, } from "./jsonl-replay.js"; import { startQaLabServer } from "./lab-server.js"; -import { runQaManualLane } from "./manual-lane.runtime.js"; +import { listLiveTransportQaAdapterFactories } from "./live-transports/cli.js"; import { loadNonYamlScenarioRefs } from "./live-transports/shared/live-transport-scenarios.js"; +import { runQaManualLane } from "./manual-lane.runtime.js"; import { runQaMultipass } from "./multipass.runtime.js"; import { DEFAULT_QA_LIVE_PROVIDER_MODE, getQaProvider } from "./providers/index.js"; import { @@ -899,10 +900,23 @@ export async function runQaSuiteCommand(opts: QaSuiteCommandOptions) { const primaryModel = normalizeQaOptionalModelRef(opts.primaryModel); const alternateModel = normalizeQaOptionalModelRef(opts.alternateModel); const channelDriver = normalizeQaSuiteChannelDriver(opts.channelDriver); - if (opts.channel?.trim() && channelDriver !== "crabline") { - throw new Error( - "--channel override is currently only supported with --channel-driver crabline.", - ); + if (opts.channel?.trim() && channelDriver !== "crabline" && channelDriver !== "live") { + throw new Error("--channel override requires --channel-driver crabline or live."); + } + const liveChannelId = channelDriver === "live" ? opts.channel?.trim() : undefined; + const liveAdapterFactories = liveChannelId ? listLiveTransportQaAdapterFactories() : undefined; + const liveAdapterFactory = liveChannelId + ? liveAdapterFactories?.find((factory) => factory.id === liveChannelId) + : undefined; + if (liveChannelId && !liveAdapterFactory) { + throw new Error(`unknown live QA adapter: ${liveChannelId}`); + } + const liveScenarioIds = + liveAdapterFactory && scenarioIds.length === 0 + ? [...(liveAdapterFactory.scenarioIds ?? [])] + : scenarioIds; + if (liveAdapterFactory && liveScenarioIds.length === 0) { + throw new Error(`live QA adapter ${liveChannelId} does not declare default scenarios`); } if (runner !== "host" && runner !== "multipass") { throw new Error(`--runner must be one of host or multipass, got "${opts.runner}".`); @@ -938,6 +952,12 @@ export async function runQaSuiteCommand(opts: QaSuiteCommandOptions) { if (runner === "multipass" && opts.cliAuthMode !== undefined) { throw new Error("--cli-auth-mode requires --runner host."); } + if (runner === "multipass" && liveChannelId) { + throw new Error("--channel-driver live with --channel requires --runner host."); + } + if (runtimePair && liveChannelId) { + throw new Error("--runtime-pair is not supported with a live QA adapter."); + } if (runner === "multipass") { rejectNonFlowScenarioIdsForMultipass(scenarioIds); const thinkingDefault = parseQaThinkingLevel("--thinking", opts.thinking); @@ -997,6 +1017,15 @@ export async function runQaSuiteCommand(opts: QaSuiteCommandOptions) { evidenceMode: opts.evidenceMode, transportId, channelDriver, + ...(liveChannelId + ? { + adapterFactories: liveAdapterFactories, + channelId: liveChannelId, + adapterOptions: { + repoRoot, + }, + } + : {}), channelDriverSelection, ...(opts.providerMode !== undefined ? { providerMode } : {}), primaryModel, @@ -1004,11 +1033,13 @@ export async function runQaSuiteCommand(opts: QaSuiteCommandOptions) { fastMode: opts.fastMode, ...(thinkingDefault ? { thinkingDefault } : {}), ...(claudeCliAuthMode ? { claudeCliAuthMode } : {}), - scenarioIds, + scenarioIds: liveChannelId ? liveScenarioIds : scenarioIds, ...(opts.enabledPluginIds !== undefined ? { enabledPluginIds: opts.enabledPluginIds } : {}), - ...(opts.concurrency !== undefined - ? { concurrency: parseQaPositiveIntegerOption("--concurrency", opts.concurrency) } - : {}), + ...(liveChannelId + ? { concurrency: 1 } + : opts.concurrency !== undefined + ? { concurrency: parseQaPositiveIntegerOption("--concurrency", opts.concurrency) } + : {}), ...(runtimePair ? { runtimePair } : {}), }), ); diff --git a/extensions/qa-lab/src/cli.ts b/extensions/qa-lab/src/cli.ts index 5bf3c6e9d227..4ea14507826b 100644 --- a/extensions/qa-lab/src/cli.ts +++ b/extensions/qa-lab/src/cli.ts @@ -470,7 +470,7 @@ export function registerQaLabCli(program: Command) { .option("--channel-driver ", "QA channel driver: qa-channel, crabline, or live") .option( "--channel ", - "Internal host QA channel override for --channel-driver; defaults to scenario/default", + "Channel id for --channel-driver crabline or live", ) .option("--provider-mode ", formatQaProviderModeHelp()) .option("--model ", "Primary provider/model ref") diff --git a/extensions/qa-lab/src/live-transports/cli.test.ts b/extensions/qa-lab/src/live-transports/cli.test.ts new file mode 100644 index 000000000000..9da116e46666 --- /dev/null +++ b/extensions/qa-lab/src/live-transports/cli.test.ts @@ -0,0 +1,75 @@ +// Qa Lab tests cover live transport CLI and adapter contribution discovery. +import { Command } from "commander"; +import type { QaRunnerCliContribution } from "openclaw/plugin-sdk/qa-runner-runtime"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const { listQaRunnerCliContributions, runSlack, runTelegram, runWhatsApp } = vi.hoisted(() => ({ + listQaRunnerCliContributions: vi.fn<() => QaRunnerCliContribution[]>(() => []), + runSlack: vi.fn(), + runTelegram: vi.fn(), + runWhatsApp: vi.fn(), +})); + +vi.mock("openclaw/plugin-sdk/qa-runner-runtime", () => ({ listQaRunnerCliContributions })); +vi.mock("./slack/cli.runtime.js", () => ({ runQaSlackCommand: runSlack })); +vi.mock("./telegram/cli.runtime.js", () => ({ runQaTelegramCommand: runTelegram })); +vi.mock("./whatsapp/cli.runtime.js", () => ({ runQaWhatsAppCommand: runWhatsApp })); + +import { + listLiveTransportQaAdapterFactories, + listLiveTransportQaCliRegistrations, +} from "./cli.js"; + +const matrixFactory = { + id: "matrix", + scenarioIds: ["channel-chat-baseline"], + matches: vi.fn(() => true), + create: vi.fn(), +}; + +describe("live transport QA contributions", () => { + beforeEach(() => { + vi.clearAllMocks(); + listQaRunnerCliContributions.mockReturnValue([ + { + pluginId: "qa-matrix", + commandName: "matrix", + status: "available", + registration: { + commandName: "matrix", + adapterFactory: matrixFactory, + register(qa) { + qa.command("matrix").action(() => undefined); + }, + }, + }, + ]); + }); + + it("discovers all four canonical live adapter factories without changing CLI ownership", () => { + expect(listLiveTransportQaAdapterFactories().map((factory) => factory.id)).toEqual([ + "telegram", + "slack", + "whatsapp", + "matrix", + ]); + }); + + it.each([ + ["telegram", runTelegram], + ["slack", runSlack], + ["whatsapp", runWhatsApp], + ] as const)("keeps the shipped %s command runner", async (commandName, runCommand) => { + const registration = listLiveTransportQaCliRegistrations().find( + (candidate) => candidate.commandName === commandName, + ); + const qa = new Command(); + registration?.register(qa); + + await qa.parseAsync(["node", "openclaw", commandName, "--scenario", `${commandName}-canary`]); + + expect(runCommand).toHaveBeenCalledWith( + expect.objectContaining({ scenarioIds: [`${commandName}-canary`] }), + ); + }); +}); diff --git a/extensions/qa-lab/src/live-transports/cli.ts b/extensions/qa-lab/src/live-transports/cli.ts index 5d48759ccc9f..3b40048e6f93 100644 --- a/extensions/qa-lab/src/live-transports/cli.ts +++ b/extensions/qa-lab/src/live-transports/cli.ts @@ -55,3 +55,9 @@ export function listLiveTransportQaCliRegistrations(): readonly LiveTransportQaC return liveRegistrations; } + +export function listLiveTransportQaAdapterFactories() { + return listLiveTransportQaCliRegistrations().flatMap((registration) => + registration.adapterFactory ? [registration.adapterFactory] : [], + ); +} diff --git a/extensions/qa-lab/src/live-transports/live-transport-adapters.test.ts b/extensions/qa-lab/src/live-transports/live-transport-adapters.test.ts new file mode 100644 index 000000000000..3fba7c4bd3f8 --- /dev/null +++ b/extensions/qa-lab/src/live-transports/live-transport-adapters.test.ts @@ -0,0 +1,68 @@ +// Qa Lab tests cover canonical live transport adapter factory routing. +import { describe, expect, it, vi } from "vitest"; +import { createQaBusState } from "../bus-state.js"; +import { createQaChannelTransport } from "../qa-channel-transport.js"; +import { createQaTransportAdapterFactoryRegistry } from "../qa-transport-registry.js"; + +const { createSlack, createTelegram, createWhatsApp } = vi.hoisted(() => ({ + createSlack: vi.fn(), + createTelegram: vi.fn(), + createWhatsApp: vi.fn(), +})); + +vi.mock("./slack/adapter.runtime.js", () => ({ createSlackQaTransportAdapter: createSlack })); +vi.mock("./telegram/adapter.runtime.js", () => ({ + createTelegramQaTransportAdapter: createTelegram, +})); +vi.mock("./whatsapp/adapter.runtime.js", () => ({ + createWhatsAppQaTransportAdapter: createWhatsApp, +})); + +import { slackQaAdapterFactory } from "./slack/cli.js"; +import { telegramQaAdapterFactory } from "./telegram/cli.js"; +import { whatsappQaAdapterFactory } from "./whatsapp/cli.js"; + +const factories = [ + telegramQaAdapterFactory, + slackQaAdapterFactory, + whatsappQaAdapterFactory, +] as const; + +describe("live transport adapter factories", () => { + it.each([ + ["telegram", createTelegram], + ["slack", createSlack], + ["whatsapp", createWhatsApp], + ] as const)( + "creates the canonical %s adapter through the shared registry", + async (channelId, create) => { + const adapterOptions = { sutAccountId: `${channelId}-sut` }; + const state = createQaBusState(); + const adapter = createQaChannelTransport(state); + create.mockResolvedValueOnce(adapter); + const registry = createQaTransportAdapterFactoryRegistry(factories); + + const created = await registry.create({ + channelId, + adapterOptions, + driver: "live", + outputDir: ".artifacts/qa-e2e", + state, + }); + + expect(created.adapter.id).toBe(adapter.id); + expect(create).toHaveBeenCalledWith( + expect.objectContaining({ + adapterOptions, + channelId, + driver: "live", + messages: expect.objectContaining({ + addInboundMessage: expect.any(Function), + addOutboundMessage: expect.any(Function), + editMessage: expect.any(Function), + }), + }), + ); + }, + ); +}); diff --git a/extensions/qa-lab/src/live-transports/slack/adapter.runtime.test.ts b/extensions/qa-lab/src/live-transports/slack/adapter.runtime.test.ts new file mode 100644 index 000000000000..b7dd468e4453 --- /dev/null +++ b/extensions/qa-lab/src/live-transports/slack/adapter.runtime.test.ts @@ -0,0 +1,42 @@ +// Qa Lab tests cover Slack live adapter message reconciliation. +import { describe, expect, it } from "vitest"; +import { createQaBusState } from "../../bus-state.js"; +import { testing } from "./adapter.runtime.js"; + +describe("Slack live adapter reconciliation", () => { + it("records streamed updates to the same Slack timestamp as bus edits", async () => { + const state = createQaBusState(); + const busMessageIds = new Map(); + const observedText = new Map(); + const messages: Parameters[0]["messages"] = { + addInboundMessage: (input) => state.addInboundMessage(input), + addOutboundMessage: (input) => state.addOutboundMessage(input), + editMessage: (input) => state.editMessage(input), + }; + const base = { + accountId: "sut", + busMessageIds, + logicalConversationId: "C123", + messages, + observedText, + sutUserId: "U123", + }; + + await testing.recordSlackObservedMessage({ + ...base, + message: { text: "QA-", ts: "123.000001", user: "U123" }, + }); + await testing.recordSlackObservedMessage({ + ...base, + message: { text: "QA-CHANNEL-BASELINE-OK", ts: "123.000001", user: "U123" }, + }); + + const snapshot = state.getSnapshot(); + expect(snapshot.messages).toHaveLength(1); + expect(snapshot.messages[0]?.text).toBe("QA-CHANNEL-BASELINE-OK"); + expect(snapshot.events.map((event) => event.kind)).toEqual([ + "outbound-message", + "message-edited", + ]); + }); +}); diff --git a/extensions/qa-lab/src/live-transports/slack/adapter.runtime.ts b/extensions/qa-lab/src/live-transports/slack/adapter.runtime.ts new file mode 100644 index 000000000000..9dc9626ff718 --- /dev/null +++ b/extensions/qa-lab/src/live-transports/slack/adapter.runtime.ts @@ -0,0 +1,193 @@ +// Qa Lab plugin module implements Slack live transport adapter behavior. +import { createSlackWebClient, createSlackWriteClient } from "@openclaw/slack/api.js"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import type { QaRunnerCliRegistration } from "openclaw/plugin-sdk/qa-runner-runtime"; +import { + acquireQaCredentialLease, + startQaCredentialLeaseHeartbeat, +} from "../shared/credential-lease.runtime.js"; +import { __testing as slackLive } from "./slack-live.runtime.js"; + +type AdapterFactory = NonNullable; +type FactoryContext = Parameters[0]; +type AdapterDefinition = Awaited>; +type SlackRuntimeEnv = ReturnType; +type SlackObservedMessage = Awaited>[number]; + +async function recordSlackObservedMessage(params: { + accountId: string; + busMessageIds: Map; + logicalConversationId: string; + message: SlackObservedMessage; + messages: FactoryContext["messages"]; + observedText: Map; + sutUserId: string; +}): Promise { + const ts = params.message.ts?.trim(); + if (!ts || params.message.user !== params.sutUserId) { + return undefined; + } + const text = params.message.text ?? ""; + if (params.observedText.get(ts) === text) { + return undefined; + } + params.observedText.set(ts, text); + const existingMessageId = params.busMessageIds.get(ts); + if (existingMessageId) { + await params.messages.editMessage({ + accountId: params.accountId, + messageId: existingMessageId, + text, + }); + return ts; + } + const outbound = await params.messages.addOutboundMessage({ + accountId: params.accountId, + to: `channel:${params.logicalConversationId}`, + senderId: params.message.user, + text, + timestamp: Number(ts.split(".")[0]) * 1_000, + threadId: params.message.thread_ts + ? params.busMessageIds.get(params.message.thread_ts) + : undefined, + }); + params.busMessageIds.set(ts, outbound.id); + return ts; +} + +export async function createSlackQaTransportAdapter( + context: FactoryContext, +): Promise { + const options = context.adapterOptions ?? {}; + const lease = await acquireQaCredentialLease({ + kind: "slack", + source: options.credentialSource, + role: options.credentialRole, + resolveEnvPayload: () => slackLive.resolveSlackQaRuntimeEnv(), + parsePayload: slackLive.parseSlackQaCredentialPayload, + }); + const heartbeat = startQaCredentialLeaseHeartbeat(lease); + const runtimeEnv = lease.payload; + let driverIdentity: Awaited>; + let sutIdentity: Awaited>; + try { + [driverIdentity, sutIdentity] = await Promise.all([ + slackLive.getSlackIdentity(runtimeEnv.driverBotToken), + slackLive.getSlackIdentity(runtimeEnv.sutBotToken), + ]); + } catch (error) { + await heartbeat.stop(); + await lease.release(); + throw error; + } + const driverClient = createSlackWriteClient(runtimeEnv.driverBotToken); + const sutClient = createSlackWebClient(runtimeEnv.sutBotToken); + const accountId = options.sutAccountId?.trim() || "sut"; + let oldestTs = `${Math.floor(Date.now() / 1_000)}.000000`; + let stopped = false; + let pollingError: Error | undefined; + let logicalConversationId = runtimeEnv.channelId; + const observedText = new Map(); + const nativeMessageIds = new Map(); + const busMessageIds = new Map(); + const polling = (async () => { + for (;;) { + if (stopped) { + return; + } + const messages = await slackLive.listSlackMessages({ + channelId: runtimeEnv.channelId, + client: sutClient, + oldestTs, + }); + for (const message of messages.toReversed()) { + const observedTs = await recordSlackObservedMessage({ + accountId, + busMessageIds, + logicalConversationId, + message, + messages: context.messages, + observedText, + sutUserId: sutIdentity.userId, + }); + if (observedTs) { + oldestTs = observedTs; + } + } + await new Promise((resolve) => { + setTimeout(resolve, 500); + }); + } + })().catch((error: unknown) => { + if (!stopped) { + pollingError = error instanceof Error ? error : new Error(String(error)); + } + }); + + return { + id: "slack", + label: "Slack live", + accountId, + requiredPluginIds: ["slack"], + supportedActions: [], + assertTransportHealthy() { + if (pollingError) { + throw pollingError; + } + heartbeat.throwIfFailed(); + }, + async sendInbound(input) { + heartbeat.throwIfFailed(); + logicalConversationId = input.conversation.id; + const text = input.text.replaceAll("@openclaw", `<@${sutIdentity.userId}>`); + const nativeThreadTs = input.threadId ? nativeMessageIds.get(input.threadId) : undefined; + const sent = await slackLive.sendSlackChannelMessage({ + channelId: runtimeEnv.channelId, + client: driverClient, + text, + threadTs: nativeThreadTs, + }); + const message = await context.messages.addInboundMessage({ + ...input, + accountId, + senderId: driverIdentity.userId, + }); + nativeMessageIds.set(message.id, sent.ts); + busMessageIds.set(sent.ts, message.id); + return message; + }, + resetTransport: () => { + logicalConversationId = runtimeEnv.channelId; + nativeMessageIds.clear(); + busMessageIds.clear(); + }, + createGatewayConfig: () => + slackLive.buildSlackQaConfig({} as OpenClawConfig, { + channelId: runtimeEnv.channelId, + driverBotUserId: driverIdentity.userId, + sutAccountId: accountId, + sutAppToken: runtimeEnv.sutAppToken, + sutBotToken: runtimeEnv.sutBotToken, + }), + waitReady: async ({ gateway }) => + await slackLive.waitForSlackChannelStable(gateway as never, accountId, "connected"), + buildAgentDelivery: () => ({ + channel: "slack", + to: `channel:${runtimeEnv.channelId}`, + replyChannel: "slack", + replyTo: `channel:${runtimeEnv.channelId}`, + }), + async handleAction() { + throw new Error("Slack live QA adapter does not implement transport actions"); + }, + createReportNotes: () => ["Runs through the Slack live adapter and shared QA suite host."], + async cleanup() { + stopped = true; + await polling.catch(() => undefined); + await heartbeat.stop(); + await lease.release(); + }, + }; +} + +export const testing = { recordSlackObservedMessage }; diff --git a/extensions/qa-lab/src/live-transports/slack/cli.ts b/extensions/qa-lab/src/live-transports/slack/cli.ts index 2935ddc54986..e9cc6d40d179 100644 --- a/extensions/qa-lab/src/live-transports/slack/cli.ts +++ b/extensions/qa-lab/src/live-transports/slack/cli.ts @@ -6,20 +6,34 @@ import { type LiveTransportQaCommandOptions, } from "../shared/live-transport-cli.js"; +type SlackQaAdapterRuntime = typeof import("./adapter.runtime.js"); type SlackQaCliRuntime = typeof import("./cli.runtime.js"); +const loadSlackQaAdapterRuntime = createLazyCliRuntimeLoader( + () => import("./adapter.runtime.js"), +); const loadSlackQaCliRuntime = createLazyCliRuntimeLoader( () => import("./cli.runtime.js"), ); async function runQaSlack(opts: LiveTransportQaCommandOptions) { - const runtime = await loadSlackQaCliRuntime(); - await runtime.runQaSlackCommand(opts); + await (await loadSlackQaCliRuntime()).runQaSlackCommand(opts); } +export const slackQaAdapterFactory: NonNullable = + { + id: "slack", + scenarioIds: ["channel-chat-baseline"], + matches: ({ channelId, driver }) => driver === "live" && channelId === "slack", + async create(context) { + return await (await loadSlackQaAdapterRuntime()).createSlackQaTransportAdapter(context); + }, + }; + export const slackQaCliRegistration: LiveTransportQaCliRegistration = createLiveTransportQaCliRegistration({ commandName: "slack", + adapterFactory: slackQaAdapterFactory, credentialOptions: { sourceDescription: "Credential source for Slack QA: env or convex (default: env)", roleDescription: @@ -27,7 +41,7 @@ export const slackQaCliRegistration: LiveTransportQaCliRegistration = }, description: "Run the Slack live QA lane against a private bot-to-bot channel harness", outputDirHelp: "Slack QA artifact directory", + run: runQaSlack, scenarioHelp: "Run only the named Slack QA scenario (repeatable)", sutAccountHelp: "Temporary Slack account id inside the QA gateway config", - run: runQaSlack, }); diff --git a/extensions/qa-lab/src/live-transports/slack/slack-live.runtime.ts b/extensions/qa-lab/src/live-transports/slack/slack-live.runtime.ts index 632794f18fdb..64f6bfd669b3 100644 --- a/extensions/qa-lab/src/live-transports/slack/slack-live.runtime.ts +++ b/extensions/qa-lab/src/live-transports/slack/slack-live.runtime.ts @@ -185,15 +185,6 @@ type SlackAuthIdentity = { userId: string; }; -type SlackMessage = { - bot_id?: string; - blocks?: unknown[]; - text?: string; - thread_ts?: string; - ts?: string; - user?: string; -}; - type SlackObservedMessage = { botId?: string; channelId: string; @@ -330,6 +321,8 @@ const slackHistoryMessageSchema = z.object({ user: z.string().optional(), }); +type SlackMessage = Omit, "ts"> & { ts?: string }; + const slackHistorySchema = z.object({ ok: z.boolean().optional(), messages: z.array(slackHistoryMessageSchema).optional(), @@ -2961,6 +2954,7 @@ export const testing = { extractSlackNativeApprovalId, findPendingCodexPluginApprovalRecord, findScenario, + getSlackIdentity, isSlackChannelReadyForQa, matchesSlackApprovalResolvedUpdate, matchesSlackApprovalPromptText, @@ -2975,9 +2969,12 @@ export const testing = { resolveApprovalDecision, resolveSlackQaSutAccountId, resolveSlackQaRuntimeEnv, + sendSlackChannelMessage, + listSlackMessages, SLACK_QA_STANDARD_SCENARIO_IDS, toSlackQaScenarioArtifactResults, waitForSlackNoReply, waitForSlackReaction, + waitForSlackChannelStable, }; export { testing as __testing }; diff --git a/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts b/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts new file mode 100644 index 000000000000..576cf4548ee1 --- /dev/null +++ b/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts @@ -0,0 +1,182 @@ +// Qa Lab plugin module implements Telegram live transport adapter behavior. +import type { TelegramBotUpdate } from "@openclaw/telegram/api.js"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import type { QaRunnerCliRegistration } from "openclaw/plugin-sdk/qa-runner-runtime"; +import { + acquireQaCredentialLease, + startQaCredentialLeaseHeartbeat, +} from "../shared/credential-lease.runtime.js"; +import { __testing as telegramLive } from "./telegram-live.runtime.js"; + +type AdapterFactory = NonNullable; +type FactoryContext = Parameters[0]; +type AdapterDefinition = Awaited>; +type TelegramRuntimeEnv = ReturnType; + +export async function createTelegramQaTransportAdapter( + context: FactoryContext, +): Promise { + const options = context.adapterOptions ?? {}; + const credentialLease = await acquireQaCredentialLease({ + kind: "telegram", + source: options.credentialSource, + role: options.credentialRole, + resolveEnvPayload: () => telegramLive.resolveTelegramQaRuntimeEnv(), + parsePayload: telegramLive.parseTelegramQaCredentialPayload, + }); + const heartbeat = startQaCredentialLeaseHeartbeat(credentialLease); + const runtimeEnv = credentialLease.payload; + let driverIdentity: { id: number; username?: string }; + let sutIdentity: { id: number; username?: string }; + let offset: number; + try { + [driverIdentity, sutIdentity, offset] = await Promise.all([ + telegramLive.callTelegramApi<{ id: number; username?: string }>( + runtimeEnv.driverToken, + "getMe", + ), + telegramLive.callTelegramApi<{ id: number; username?: string }>(runtimeEnv.sutToken, "getMe"), + telegramLive.flushTelegramUpdates(runtimeEnv.driverToken), + ]); + } catch (error) { + await heartbeat.stop(); + await credentialLease.release(); + throw error; + } + let stopped = false; + let pollingError: Error | undefined; + let logicalConversationId = runtimeEnv.groupId; + let logicalConversationKind: "channel" | "direct" | "group" = "channel"; + const nativeMessageIds = new Map(); + const busMessageIds = new Map(); + const poll = async () => { + for (;;) { + if (stopped) { + return; + } + const updates = await telegramLive.callTelegramApi( + runtimeEnv.driverToken, + "getUpdates", + { offset, timeout: 1, allowed_updates: ["message", "edited_message"] }, + 6_000, + ); + for (const update of updates) { + offset = Math.max(offset, update.update_id + 1); + const message = update.edited_message ?? update.message; + if (!message?.from?.id || message.from.id !== sutIdentity.id) { + continue; + } + const existingMessageId = busMessageIds.get(message.message_id); + if (update.edited_message && existingMessageId) { + await context.messages.editMessage({ + accountId: options.sutAccountId?.trim() || "sut", + messageId: existingMessageId, + text: message.text ?? message.caption ?? "", + }); + continue; + } + const outbound = await context.messages.addOutboundMessage({ + accountId: options.sutAccountId?.trim() || "sut", + to: `${logicalConversationKind}:${logicalConversationId}`, + senderId: String(message.from.id), + senderName: message.from.username, + text: message.text ?? message.caption ?? "", + timestamp: message.date * 1_000, + replyToId: message.reply_to_message?.message_id + ? busMessageIds.get(message.reply_to_message.message_id) + : undefined, + }); + busMessageIds.set(message.message_id, outbound.id); + } + } + }; + const polling = poll().catch((error: unknown) => { + if (!stopped) { + pollingError = error instanceof Error ? error : new Error(String(error)); + } + }); + const accountId = options.sutAccountId?.trim() || "sut"; + + return { + id: "telegram", + label: "Telegram live", + accountId, + requiredPluginIds: ["telegram"], + supportedActions: [], + assertTransportHealthy() { + if (pollingError) { + throw pollingError; + } + heartbeat.throwIfFailed(); + }, + async sendInbound(input) { + heartbeat.throwIfFailed(); + logicalConversationId = input.conversation.id; + logicalConversationKind = input.conversation.kind; + const text = sutIdentity.username + ? input.text.replaceAll("@openclaw", `@${sutIdentity.username}`) + : input.text; + const nativeReplyToId = input.replyToId ? nativeMessageIds.get(input.replyToId) : undefined; + const sent = await telegramLive.callTelegramApi<{ message_id: number }>( + runtimeEnv.driverToken, + "sendMessage", + { + chat_id: runtimeEnv.groupId, + text, + disable_notification: true, + ...(nativeReplyToId + ? { + reply_parameters: { + message_id: nativeReplyToId, + allow_sending_without_reply: true, + }, + } + : {}), + }, + ); + const message = await context.messages.addInboundMessage({ + ...input, + accountId, + senderId: String(driverIdentity.id), + senderName: driverIdentity.username, + }); + nativeMessageIds.set(message.id, sent.message_id); + busMessageIds.set(sent.message_id, message.id); + return message; + }, + resetTransport: () => { + logicalConversationId = runtimeEnv.groupId; + logicalConversationKind = "channel"; + nativeMessageIds.clear(); + busMessageIds.clear(); + }, + createGatewayConfig: () => + telegramLive.buildTelegramQaConfig({} as OpenClawConfig, { + groupId: runtimeEnv.groupId, + sutToken: runtimeEnv.sutToken, + driverBotId: driverIdentity.id, + sutAccountId: accountId, + }), + waitReady: async ({ gateway, timeoutMs, pollIntervalMs }) => + await telegramLive.waitForTelegramChannelRunning(gateway as never, accountId, { + timeoutMs, + pollMs: pollIntervalMs, + }), + buildAgentDelivery: () => ({ + channel: "telegram", + to: runtimeEnv.groupId, + replyChannel: "telegram", + replyTo: runtimeEnv.groupId, + }), + async handleAction() { + throw new Error("Telegram live QA adapter does not implement transport actions"); + }, + createReportNotes: () => ["Runs through the Telegram live adapter and shared QA suite host."], + async cleanup() { + stopped = true; + await polling.catch(() => undefined); + await heartbeat.stop(); + await credentialLease.release(); + }, + }; +} diff --git a/extensions/qa-lab/src/live-transports/telegram/cli.ts b/extensions/qa-lab/src/live-transports/telegram/cli.ts index 550901fb7178..d20ff2615a49 100644 --- a/extensions/qa-lab/src/live-transports/telegram/cli.ts +++ b/extensions/qa-lab/src/live-transports/telegram/cli.ts @@ -6,20 +6,35 @@ import { type LiveTransportQaCommandOptions, } from "../shared/live-transport-cli.js"; +type TelegramQaAdapterRuntime = typeof import("./adapter.runtime.js"); type TelegramQaCliRuntime = typeof import("./cli.runtime.js"); +const loadTelegramQaAdapterRuntime = createLazyCliRuntimeLoader( + () => import("./adapter.runtime.js"), +); const loadTelegramQaCliRuntime = createLazyCliRuntimeLoader( () => import("./cli.runtime.js"), ); async function runQaTelegram(opts: LiveTransportQaCommandOptions) { - const runtime = await loadTelegramQaCliRuntime(); - await runtime.runQaTelegramCommand(opts); + await (await loadTelegramQaCliRuntime()).runQaTelegramCommand(opts); } +export const telegramQaAdapterFactory: NonNullable< + LiveTransportQaCliRegistration["adapterFactory"] +> = { + id: "telegram", + scenarioIds: ["channel-chat-baseline"], + matches: ({ channelId, driver }) => driver === "live" && channelId === "telegram", + async create(context) { + return await (await loadTelegramQaAdapterRuntime()).createTelegramQaTransportAdapter(context); + }, +}; + export const telegramQaCliRegistration: LiveTransportQaCliRegistration = createLiveTransportQaCliRegistration({ commandName: "telegram", + adapterFactory: telegramQaAdapterFactory, credentialOptions: { sourceDescription: "Credential source for Telegram QA: env or convex (default: env)", roleDescription: @@ -28,7 +43,7 @@ export const telegramQaCliRegistration: LiveTransportQaCliRegistration = description: "Run the manual Telegram live QA lane against a private bot-to-bot group harness", listScenariosHelp: "Print available Telegram scenario ids and exit", outputDirHelp: "Telegram QA artifact directory", + run: runQaTelegram, scenarioHelp: "Run only the named Telegram QA scenario (repeatable)", sutAccountHelp: "Temporary Telegram account id inside the QA gateway config", - run: runQaTelegram, }); diff --git a/extensions/qa-lab/src/live-transports/telegram/telegram-live.runtime.ts b/extensions/qa-lab/src/live-transports/telegram/telegram-live.runtime.ts index 96f730de7609..a190c24dab49 100644 --- a/extensions/qa-lab/src/live-transports/telegram/telegram-live.runtime.ts +++ b/extensions/qa-lab/src/live-transports/telegram/telegram-live.runtime.ts @@ -2,6 +2,7 @@ import { randomUUID } from "node:crypto"; import fs from "node:fs/promises"; import path from "node:path"; +import type { TelegramBotMessage, TelegramBotUpdate } from "@openclaw/telegram/api.js"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime"; import { @@ -218,32 +219,22 @@ type TelegramRichMessage = { blocks?: unknown[]; }; -type TelegramMessage = { - message_id: number; - date: number; - text?: string; - caption?: string; +type TelegramMessage = Pick & + Partial> & { + audio?: unknown; + chat: { id: number }; + document?: unknown; + from?: Pick, "id" | "is_bot" | "username">; + photo?: unknown[]; rich_message?: TelegramRichMessage; reply_markup?: TelegramReplyMarkup; reply_to_message?: { message_id?: number }; - from?: { - id?: number; - is_bot?: boolean; - username?: string; - }; - chat: { - id: number; - }; - photo?: unknown[]; - document?: unknown; - audio?: unknown; + sticker?: unknown; video?: unknown; voice?: unknown; - sticker?: unknown; }; -type TelegramUpdate = { - update_id: number; +type TelegramUpdate = Pick & { edited_message?: TelegramMessage; message?: TelegramMessage; }; @@ -841,7 +832,7 @@ function normalizeTelegramObservedMessage(update: TelegramUpdate): TelegramObser messageId: message.message_id, chatId: message.chat.id, senderId: message.from.id, - senderIsBot: message.from.is_bot === true, + senderIsBot: message.from.is_bot, senderUsername: message.from.username, text: selectTelegramObservedText(message), caption: message.caption, @@ -2258,6 +2249,7 @@ export const testing = { assertTelegramScenarioReply, classifyCanaryReply, findScenario, + flushTelegramUpdates, isTelegramObservedMessageTimeoutError, listTelegramQaScenarioCatalog, matchesTelegramScenarioReply, diff --git a/extensions/qa-lab/src/live-transports/whatsapp/adapter.runtime.ts b/extensions/qa-lab/src/live-transports/whatsapp/adapter.runtime.ts new file mode 100644 index 000000000000..b997cb4de01b --- /dev/null +++ b/extensions/qa-lab/src/live-transports/whatsapp/adapter.runtime.ts @@ -0,0 +1,182 @@ +// Qa Lab plugin module implements WhatsApp live transport adapter behavior. +import fs from "node:fs/promises"; +import path from "node:path"; +import { startWhatsAppQaDriverSession } from "@openclaw/whatsapp/api.js"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import type { QaRunnerCliRegistration } from "openclaw/plugin-sdk/qa-runner-runtime"; +import { resolvePreferredOpenClawTmpDir } from "openclaw/plugin-sdk/temp-path"; +import { + acquireQaCredentialLease, + startQaCredentialLeaseHeartbeat, +} from "../shared/credential-lease.runtime.js"; +import { __testing as whatsappLive } from "./whatsapp-live.runtime.js"; + +type AdapterFactory = NonNullable; +type FactoryContext = Parameters[0]; +type AdapterDefinition = Awaited>; +type WhatsAppRuntimeEnv = ReturnType; + +export async function createWhatsAppQaTransportAdapter( + context: FactoryContext, +): Promise { + const options = context.adapterOptions ?? {}; + const lease = await acquireQaCredentialLease({ + kind: "whatsapp", + source: options.credentialSource, + role: options.credentialRole, + resolveEnvPayload: () => whatsappLive.resolveWhatsAppQaRuntimeEnv(), + parsePayload: whatsappLive.parseWhatsAppQaCredentialPayload, + }); + const heartbeat = startQaCredentialLeaseHeartbeat(lease); + const runtimeEnv = lease.payload; + let authRoot: string | undefined; + let driver: Awaited> | undefined; + let sutAuthDir: string; + try { + authRoot = await fs.mkdtemp( + path.join(resolvePreferredOpenClawTmpDir(), "openclaw-whatsapp-qa-adapter-"), + ); + const [driverAuthDir, unpackedSutAuthDir] = await Promise.all([ + whatsappLive.unpackWhatsAppAuthArchive({ + archiveBase64: runtimeEnv.driverAuthArchiveBase64, + clearSignalSessions: true, + label: "driver-auth", + parentDir: authRoot, + }), + whatsappLive.unpackWhatsAppAuthArchive({ + archiveBase64: runtimeEnv.sutAuthArchiveBase64, + clearSignalSessions: true, + label: "sut-auth", + parentDir: authRoot, + }), + ]); + sutAuthDir = unpackedSutAuthDir; + driver = await startWhatsAppQaDriverSession({ authDir: driverAuthDir }); + } catch (error) { + await driver?.close().catch(() => undefined); + await heartbeat.stop(); + await lease.release(); + if (authRoot) { + await fs.rm(authRoot, { force: true, recursive: true }); + } + throw error; + } + const accountId = options.sutAccountId?.trim() || "sut"; + const targets = whatsappLive.resolveWhatsAppQaMessageTargets({ + driverPhoneE164: runtimeEnv.driverPhoneE164, + scenarioTarget: "dm", + sutPhoneE164: runtimeEnv.sutPhoneE164, + }); + let observedCount = driver.getObservedMessages().length; + let stopped = false; + let pollingError: Error | undefined; + let logicalConversationId = targets.gatewayTarget; + const nativeMessageIds = new Map(); + const busMessageIds = new Map(); + const polling = (async () => { + for (;;) { + if (stopped) { + return; + } + const messages = driver.getObservedMessages(); + for (const message of messages.slice(observedCount)) { + observedCount += 1; + if (message.fromPhoneE164 !== runtimeEnv.sutPhoneE164) { + continue; + } + await context.messages.addOutboundMessage({ + accountId, + to: `dm:${logicalConversationId}`, + senderId: message.fromPhoneE164, + text: message.text, + timestamp: Date.parse(message.observedAt), + replyToId: message.quoted?.messageId + ? busMessageIds.get(message.quoted.messageId) + : undefined, + }); + } + await new Promise((resolve) => { + setTimeout(resolve, 500); + }); + } + })().catch((error: unknown) => { + if (!stopped) { + pollingError = error instanceof Error ? error : new Error(String(error)); + } + }); + + return { + id: "whatsapp", + label: "WhatsApp live", + accountId, + requiredPluginIds: ["whatsapp"], + supportedActions: [], + assertTransportHealthy() { + if (pollingError) { + throw pollingError; + } + heartbeat.throwIfFailed(); + }, + async sendInbound(input) { + heartbeat.throwIfFailed(); + logicalConversationId = input.conversation.id; + const quotedMessageId = input.replyToId ? nativeMessageIds.get(input.replyToId) : undefined; + const sent = await driver.sendText( + targets.driverTarget, + input.text, + quotedMessageId + ? { + quotedMessageKey: { + id: quotedMessageId, + remoteJid: targets.driverTarget, + fromMe: true, + }, + } + : undefined, + ); + const message = await context.messages.addInboundMessage({ + ...input, + accountId, + senderId: runtimeEnv.driverPhoneE164, + }); + if (sent.messageId) { + nativeMessageIds.set(message.id, sent.messageId); + busMessageIds.set(sent.messageId, message.id); + } + return message; + }, + resetTransport: () => { + logicalConversationId = targets.gatewayTarget; + nativeMessageIds.clear(); + busMessageIds.clear(); + }, + createGatewayConfig: () => + whatsappLive.buildWhatsAppQaConfig({} as OpenClawConfig, { + allowFrom: [runtimeEnv.driverPhoneE164], + authDir: sutAuthDir, + dmPolicy: "allowlist", + groupJid: runtimeEnv.groupJid, + sutAccountId: accountId, + }), + waitReady: async ({ gateway }) => + await whatsappLive.waitForWhatsAppChannelStable(gateway as never, accountId), + buildAgentDelivery: () => ({ + channel: "whatsapp", + to: targets.gatewayTarget, + replyChannel: "whatsapp", + replyTo: targets.gatewayTarget, + }), + async handleAction() { + throw new Error("WhatsApp live QA adapter does not implement transport actions"); + }, + createReportNotes: () => ["Uses the WhatsApp live adapter."], + async cleanup() { + stopped = true; + await polling.catch(() => undefined); + await driver.close(); + await heartbeat.stop(); + await lease.release(); + await fs.rm(authRoot, { force: true, recursive: true }); + }, + }; +} diff --git a/extensions/qa-lab/src/live-transports/whatsapp/cli.ts b/extensions/qa-lab/src/live-transports/whatsapp/cli.ts index 9f13a7627149..a28b4cf4a93f 100644 --- a/extensions/qa-lab/src/live-transports/whatsapp/cli.ts +++ b/extensions/qa-lab/src/live-transports/whatsapp/cli.ts @@ -6,20 +6,35 @@ import { type LiveTransportQaCommandOptions, } from "../shared/live-transport-cli.js"; +type WhatsAppQaAdapterRuntime = typeof import("./adapter.runtime.js"); type WhatsAppQaCliRuntime = typeof import("./cli.runtime.js"); +const loadWhatsAppQaAdapterRuntime = createLazyCliRuntimeLoader( + () => import("./adapter.runtime.js"), +); const loadWhatsAppQaCliRuntime = createLazyCliRuntimeLoader( () => import("./cli.runtime.js"), ); async function runQaWhatsApp(opts: LiveTransportQaCommandOptions) { - const runtime = await loadWhatsAppQaCliRuntime(); - await runtime.runQaWhatsAppCommand(opts); + await (await loadWhatsAppQaCliRuntime()).runQaWhatsAppCommand(opts); } +export const whatsappQaAdapterFactory: NonNullable< + LiveTransportQaCliRegistration["adapterFactory"] +> = { + id: "whatsapp", + scenarioIds: ["dm-chat-baseline"], + matches: ({ channelId, driver }) => driver === "live" && channelId === "whatsapp", + async create(context) { + return await (await loadWhatsAppQaAdapterRuntime()).createWhatsAppQaTransportAdapter(context); + }, +}; + export const whatsappQaCliRegistration: LiveTransportQaCliRegistration = createLiveTransportQaCliRegistration({ commandName: "whatsapp", + adapterFactory: whatsappQaAdapterFactory, credentialOptions: { sourceDescription: "Credential source for WhatsApp QA: env or convex (default: env)", roleDescription: @@ -27,7 +42,7 @@ export const whatsappQaCliRegistration: LiveTransportQaCliRegistration = }, description: "Run the WhatsApp live QA lane against two pre-linked Web sessions", outputDirHelp: "WhatsApp QA artifact directory", + run: runQaWhatsApp, scenarioHelp: "Run only the named WhatsApp QA scenario (repeatable)", sutAccountHelp: "Temporary WhatsApp account id inside the QA gateway config", - run: runQaWhatsApp, }); diff --git a/extensions/qa-lab/src/live-transports/whatsapp/whatsapp-live.runtime.ts b/extensions/qa-lab/src/live-transports/whatsapp/whatsapp-live.runtime.ts index 1c978461c912..2de1935586a4 100644 --- a/extensions/qa-lab/src/live-transports/whatsapp/whatsapp-live.runtime.ts +++ b/extensions/qa-lab/src/live-transports/whatsapp/whatsapp-live.runtime.ts @@ -4725,6 +4725,7 @@ export const testing = { runWhatsAppApprovalScenario, runWhatsAppStructuredInboundChecks, waitForScenarioObservedMessage, + waitForWhatsAppChannelStable, redactWhatsAppQaScenarioResults, resolveWhatsAppQaMessageTargets, resolveWhatsAppQaRuntimeEnv, diff --git a/extensions/qa-lab/src/qa-transport-registry.test.ts b/extensions/qa-lab/src/qa-transport-registry.test.ts index 79dc5f7be1e9..9430456d374b 100644 --- a/extensions/qa-lab/src/qa-transport-registry.test.ts +++ b/extensions/qa-lab/src/qa-transport-registry.test.ts @@ -1,7 +1,6 @@ // Qa Lab tests cover qa transport registry plugin behavior. import { describe, expect, it, vi } from "vitest"; import { createQaBusState } from "./bus-state.js"; -import { createQaChannelTransport } from "./qa-channel-transport.js"; import { createQaTransportAdapter, createQaTransportAdapterFactoryRegistry, @@ -11,6 +10,31 @@ import { } from "./qa-transport-registry.js"; import type { QaTransportAdapter } from "./qa-transport.js"; +function createAdapterDefinition(cleanup?: () => Promise) { + const state = createQaBusState(); + return { + id: "selected", + label: "Selected", + accountId: "sut", + requiredPluginIds: [], + supportedActions: [], + async sendInbound(input: Parameters[0]) { + return state.addInboundMessage(input); + }, + createGatewayConfig: () => ({}), + async waitReady() {}, + buildAgentDelivery: ({ target }: { target: string }) => ({ + channel: "selected", + to: target, + replyChannel: "selected", + replyTo: target, + }), + async handleAction() {}, + createReportNotes: () => [], + ...(cleanup ? { cleanup } : {}), + }; +} + function createFactoryContext( overrides: Partial = {}, ): QaTransportFactoryContext { @@ -39,35 +63,42 @@ describe("qa transport registry", () => { }); it("selects an injected matching factory", async () => { - const adapter = createQaChannelTransport(createQaBusState()); - const skippedCreate = vi.fn(async () => adapter); - const selectedCreate = vi.fn(async () => adapter); + const definition = createAdapterDefinition(); + const skippedCreate = vi.fn(async () => definition); + const selectedCreate = vi.fn(async () => definition); const factories: QaTransportAdapterFactory[] = [ { id: "skipped", matches: () => false, create: skippedCreate }, { id: "selected", matches: () => true, create: selectedCreate }, ]; const registry = createQaTransportAdapterFactoryRegistry(factories); - const created = await registry.create(createFactoryContext()); + const created = await registry.create( + createFactoryContext({ channelId: "selected", driver: "live" }), + ); - expect(created.adapter).toBe(adapter); + expect(created.adapter).toMatchObject({ + id: definition.id, + label: definition.label, + state: expect.any(Object), + }); expect(skippedCreate).not.toHaveBeenCalled(); expect(selectedCreate).toHaveBeenCalledOnce(); }); it("returns cleanup owned by the selected adapter", async () => { const cleanup = vi.fn(async () => undefined); - const adapter: QaTransportAdapter = createQaChannelTransport(createQaBusState()); - adapter.cleanup = cleanup; + const definition = createAdapterDefinition(cleanup); const factory: QaTransportAdapterFactory = { id: "cleanup", matches: () => true, async create() { - return adapter; + return definition; }, }; const registry = createQaTransportAdapterFactoryRegistry([factory]); - const created = await registry.create(createFactoryContext()); + const created = await registry.create( + createFactoryContext({ channelId: "cleanup", driver: "live" }), + ); await created.cleanup(); @@ -75,10 +106,10 @@ describe("qa transport registry", () => { }); it("reports no-match and startup failures with transport context", async () => { - const context = createFactoryContext(); + const context = createFactoryContext({ channelId: "missing", driver: "live" }); const emptyRegistry = createQaTransportAdapterFactoryRegistry([]); await expect(emptyRegistry.create(context)).rejects.toThrow( - "no QA transport factory for qa-channel:qa-channel", + "no QA transport factory for live:missing", ); const brokenRegistry = createQaTransportAdapterFactoryRegistry([ @@ -91,7 +122,7 @@ describe("qa transport registry", () => { }, ]); await expect(brokenRegistry.create(context)).rejects.toThrow( - "broken failed to create QA transport qa-channel:qa-channel: provider boot failed", + "failed to create QA transport live:missing: provider boot failed", ); }); }); diff --git a/extensions/qa-lab/src/qa-transport-registry.ts b/extensions/qa-lab/src/qa-transport-registry.ts index 7ceadf647d44..df8dc891c4ef 100644 --- a/extensions/qa-lab/src/qa-transport-registry.ts +++ b/extensions/qa-lab/src/qa-transport-registry.ts @@ -1,3 +1,4 @@ +import type { QaRunnerCliRegistration } from "openclaw/plugin-sdk/qa-runner-runtime"; // Qa Lab plugin module implements qa transport registry behavior. import type { QaBusState } from "./bus-state.js"; import { @@ -5,27 +6,29 @@ import { QA_CHANNEL_DEFAULT_SUITE_CONCURRENCY, } from "./qa-channel-transport.js"; import type { QaTransportAdapter } from "./qa-transport.js"; +import { createQaStateBackedTransportAdapter } from "./qa-transport.js"; export type QaTransportId = "qa-channel"; -export type QaTransportDriver = QaTransportId | "crabline"; +export type QaTransportDriver = QaTransportId | "crabline" | "live"; export type QaTransportFactoryContext = { + adapterOptions?: Parameters< + NonNullable["create"] + >[0]["adapterOptions"]; channelId: string; driver: QaTransportDriver; outputDir: string; state: QaBusState; }; -export type QaTransportAdapterFactoryResult = { - adapter: QaTransportAdapter; +export type QaTransportAdapterFactoryResult< + TAdapter extends QaTransportAdapter = QaTransportAdapter, +> = { + adapter: TAdapter; cleanup: () => Promise; }; -export type QaTransportAdapterFactory = { - id: string; - matches: (context: Pick) => boolean; - create: (context: QaTransportFactoryContext) => Promise; -}; +export type QaTransportAdapterFactory = NonNullable; export type QaTransportAdapterFactoryRegistry = { create: (context: QaTransportFactoryContext) => Promise; @@ -33,18 +36,13 @@ export type QaTransportAdapterFactoryRegistry = { const DEFAULT_QA_TRANSPORT_ID: QaTransportId = "qa-channel"; -const QA_CHANNEL_TRANSPORT_FACTORY: QaTransportAdapterFactory = { - id: "qa-channel", - matches: ({ channelId, driver }) => driver === "qa-channel" && channelId === "qa-channel", - async create(context) { +async function createBuiltInQaTransport( + context: QaTransportFactoryContext, +): Promise { + if (context.driver === "qa-channel" && context.channelId === "qa-channel") { return createQaChannelTransport(context.state); - }, -}; - -const CRABLINE_TRANSPORT_FACTORY: QaTransportAdapterFactory = { - id: "crabline", - matches: ({ driver }) => driver === "crabline", - async create(context) { + } + if (context.driver === "crabline") { const { resolveOpenClawCrablineChannelDriverSelection } = await import("@openclaw/crabline"); const selection = resolveOpenClawCrablineChannelDriverSelection({ channel: context.channelId }); const { createQaCrablineTransportAdapter } = await import("./crabline-transport.js"); @@ -53,13 +51,9 @@ const CRABLINE_TRANSPORT_FACTORY: QaTransportAdapterFactory = { selection, state: context.state, }); - }, -}; - -const DEFAULT_QA_TRANSPORT_FACTORIES = [ - QA_CHANNEL_TRANSPORT_FACTORY, - CRABLINE_TRANSPORT_FACTORY, -] as const; + } + return undefined; +} function requireQaTransportFactory( factories: readonly QaTransportAdapterFactory[], @@ -73,20 +67,35 @@ function requireQaTransportFactory( } export function createQaTransportAdapterFactoryRegistry( - factories: readonly QaTransportAdapterFactory[] = DEFAULT_QA_TRANSPORT_FACTORIES, + factories: readonly QaTransportAdapterFactory[] = [], ): QaTransportAdapterFactoryRegistry { return { async create(context) { - const factory = requireQaTransportFactory(factories, context); let adapter: QaTransportAdapter; try { - adapter = await factory.create(context); + const builtIn = await createBuiltInQaTransport(context); + if (builtIn) { + adapter = builtIn; + } else { + const factory = requireQaTransportFactory(factories, context); + const definition = await factory.create({ + adapterOptions: context.adapterOptions, + channelId: context.channelId, + driver: context.driver, + messages: { + addInboundMessage: (input) => context.state.addInboundMessage(input), + addOutboundMessage: (input) => context.state.addOutboundMessage(input), + editMessage: (input) => context.state.editMessage(input), + }, + outputDir: context.outputDir, + }); + adapter = createQaStateBackedTransportAdapter(context.state, definition); + } } catch (error) { const message = error instanceof Error ? error.message : String(error); - throw new Error( - `${factory.id} failed to create QA transport ${context.driver}:${context.channelId}: ${message}`, - { cause: error }, - ); + throw new Error(`failed to create QA transport ${context.driver}:${context.channelId}: ${message}`, { + cause: error, + }); } return { adapter, @@ -110,8 +119,13 @@ export function normalizeQaTransportId(input?: string | null): QaTransportId { export async function createQaTransportAdapter( context: QaTransportFactoryContext, + factories?: readonly QaTransportAdapterFactory[], ): Promise { - return await qaTransportAdapterFactoryRegistry.create(context); + return await ( + factories + ? createQaTransportAdapterFactoryRegistry(factories) + : qaTransportAdapterFactoryRegistry + ).create(context); } export function defaultQaSuiteConcurrencyForTransport(id: QaTransportId): number { diff --git a/extensions/qa-lab/src/qa-transport.test.ts b/extensions/qa-lab/src/qa-transport.test.ts index cf3079fa08d8..b38c73b17013 100644 --- a/extensions/qa-lab/src/qa-transport.test.ts +++ b/extensions/qa-lab/src/qa-transport.test.ts @@ -1,7 +1,48 @@ // Qa Lab tests cover shared transport behavior. -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { createQaBusState } from "./bus-state.js"; -import { waitForQaTransportOutboundSequence } from "./qa-transport.js"; +import { + createQaStateBackedTransportAdapter, + waitForQaTransportOutboundSequence, +} from "./qa-transport.js"; + +describe("createQaStateBackedTransportAdapter", () => { + it("runs transport reset before clearing shared state", async () => { + const state = createQaBusState(); + state.addInboundMessage({ + conversation: { id: "alice", kind: "direct" }, + senderId: "alice", + text: "hello", + }); + const resetTransport = vi.fn(() => { + expect(state.getSnapshot().messages).toHaveLength(1); + }); + const adapter = createQaStateBackedTransportAdapter(state, { + id: "live", + label: "Live", + accountId: "sut", + requiredPluginIds: [], + supportedActions: [], + resetTransport, + sendInbound: async (input) => state.addInboundMessage(input), + createGatewayConfig: () => ({}), + waitReady: async () => undefined, + buildAgentDelivery: ({ target }) => ({ + channel: "live", + to: target, + replyChannel: "live", + replyTo: target, + }), + handleAction: async () => undefined, + createReportNotes: () => [], + }); + + await adapter.reset(); + + expect(resetTransport).toHaveBeenCalledOnce(); + expect(state.getSnapshot().messages).toHaveLength(0); + }); +}); describe("waitForQaTransportOutboundSequence", () => { it("returns preview and final edit events for one threaded message", async () => { diff --git a/extensions/qa-lab/src/qa-transport.ts b/extensions/qa-lab/src/qa-transport.ts index d732bbfb7de8..08fa953ba4a2 100644 --- a/extensions/qa-lab/src/qa-transport.ts +++ b/extensions/qa-lab/src/qa-transport.ts @@ -2,9 +2,11 @@ import { setTimeout as sleep } from "node:timers/promises"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { resolveTimerTimeoutMs } from "openclaw/plugin-sdk/number-runtime"; +import type { QaRunnerCliRegistration } from "openclaw/plugin-sdk/qa-runner-runtime"; import type { QaProviderMode } from "./model-selection.js"; import { extractQaFailureReplyText } from "./reply-failure.js"; import type { + QaBusEditMessageInput, QaBusEvent, QaBusInboundMessageInput, QaBusMessage, @@ -43,6 +45,7 @@ export type QaTransportState = { getSnapshot: () => QaBusStateSnapshot; addInboundMessage: (input: QaBusInboundMessageInput) => QaBusMessage | Promise; addOutboundMessage: (input: QaBusOutboundMessageInput) => QaBusMessage | Promise; + editMessage?: (input: QaBusEditMessageInput) => QaBusMessage | Promise; readMessage: ( input: QaBusReadMessageInput, ) => QaBusMessage | null | undefined | Promise; @@ -174,47 +177,23 @@ export function createFailureAwareTransportWaitForCondition(state: QaTransportSt }; } -export type QaTransportAdapter = { - id: string; - label: string; - accountId: string; - requiredPluginIds: readonly string[]; - supportedActions: readonly QaTransportActionName[]; +type QaTransportAdapterDefinition = Awaited< + ReturnType["create"]> +>; + +export type QaTransportAdapter = Omit< + QaTransportAdapterDefinition, + "assertTransportHealthy" | "resetTransport" +> & { state: QaTransportState; reset: () => Promise; - sendInbound: (input: QaBusInboundMessageInput) => Promise; - sendNativeCommand?: (input: QaTransportNativeCommandInput) => Promise; waitForNoOutbound: (input?: QaTransportWaitForNoOutboundInput) => Promise; waitForOutbound: (input: QaTransportOutboundMatch) => Promise; - waitForOutboundSequence?: ( - input: QaTransportOutboundSequenceMatch, - ) => Promise; waitForCondition: ( check: () => T | Promise | null | undefined, timeoutMs?: number, intervalMs?: number, ) => Promise; - createGatewayConfig: (params: { baseUrl: string }) => QaTransportGatewayConfig; - waitReady: (params: { - gateway: QaTransportGatewayClient; - timeoutMs?: number; - pollIntervalMs?: number; - }) => Promise; - buildAgentDelivery: (params: { target: string }) => { - channel: string; - to?: string; - replyChannel: string; - replyTo: string; - }; - createRuntimeEnvPatch?: () => NodeJS.ProcessEnv; - handleAction: (params: { - action: QaTransportActionName; - args: Record; - cfg: OpenClawConfig; - accountId?: string | null; - }) => Promise; - createReportNotes: (params: QaTransportReportParams) => string[]; - cleanup?: () => Promise; }; export abstract class QaStateBackedTransportAdapter implements QaTransportAdapter { @@ -225,14 +204,16 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte readonly supportedActions: readonly QaTransportActionName[]; readonly state: QaTransportState; readonly waitForCondition: QaTransportAdapter["waitForCondition"]; + private readonly assertTransportHealthy: () => void; - protected constructor(params: { + constructor(params: { id: string; label: string; accountId: string; requiredPluginIds: readonly string[]; supportedActions?: readonly QaTransportActionName[]; state: QaTransportState; + assertTransportHealthy?: () => void; }) { this.id = params.id; this.label = params.label; @@ -240,7 +221,17 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte this.requiredPluginIds = params.requiredPluginIds; this.supportedActions = params.supportedActions ?? []; this.state = params.state; - this.waitForCondition = createFailureAwareTransportWaitForCondition(this.state); + this.assertTransportHealthy = params.assertTransportHealthy ?? (() => undefined); + const waitForCondition = createFailureAwareTransportWaitForCondition(this.state); + this.waitForCondition = async (check, timeoutMs, intervalMs) => + await waitForCondition( + async () => { + this.assertTransportHealthy(); + return await check(); + }, + timeoutMs, + intervalMs, + ); } abstract createGatewayConfig: (params: { baseUrl: string }) => QaTransportGatewayConfig; @@ -264,6 +255,7 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte abstract createReportNotes: (params: QaTransportReportParams) => string[]; async reset() { + this.assertTransportHealthy(); await this.state.reset(); } @@ -272,8 +264,10 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte } async waitForNoOutbound(input: QaTransportWaitForNoOutboundInput = {}) { + this.assertTransportHealthy(); const quietMs = resolveTimerTimeoutMs(input.quietMs, 1_200, 0); await sleep(quietMs); + this.assertTransportHealthy(); assertNoFailureReplies(this.state, { sinceIndex: input.sinceIndex, cursorSpace: "outbound", @@ -287,6 +281,7 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte async waitForOutbound(input: QaTransportOutboundMatch) { return await waitForQaTransportCondition(() => { + this.assertTransportHealthy(); assertNoFailureReplies(this.state, { sinceIndex: input.sinceIndex, cursorSpace: "outbound", @@ -317,6 +312,52 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte } } +export function createQaStateBackedTransportAdapter( + state: QaTransportState, + params: QaTransportAdapterDefinition, +): QaTransportAdapter { + const adapter = new (class extends QaStateBackedTransportAdapter { + createGatewayConfig = params.createGatewayConfig; + waitReady = params.waitReady; + buildAgentDelivery = params.buildAgentDelivery; + handleAction = params.handleAction; + createReportNotes = params.createReportNotes; + + override sendInbound = params.sendInbound; + + override async reset() { + await params.resetTransport?.(); + await super.reset(); + } + })({ + id: params.id, + label: params.label, + accountId: params.accountId, + requiredPluginIds: params.requiredPluginIds, + supportedActions: params.supportedActions, + state, + assertTransportHealthy: params.assertTransportHealthy, + }); + Object.assign(adapter, { + ...(params.sendNativeCommand ? { sendNativeCommand: params.sendNativeCommand } : {}), + waitForOutboundSequence: + params.waitForOutboundSequence ?? + (async (input: QaTransportOutboundSequenceMatch) => + await waitForQaTransportOutboundSequence({ + input, + readEvents: () => { + params.assertTransportHealthy?.(); + return state.getSnapshot().events; + }, + })), + ...(params.createRuntimeEnvPatch + ? { createRuntimeEnvPatch: params.createRuntimeEnvPatch } + : {}), + ...(params.cleanup ? { cleanup: params.cleanup } : {}), + }); + return adapter; +} + function normalizeQaBusOutboundEvent(event: QaBusEvent): QaTransportOutboundEvent | null { switch (event.kind) { case "outbound-message": diff --git a/extensions/qa-lab/src/suite.test.ts b/extensions/qa-lab/src/suite.test.ts index 9347b05b798b..5b5433c5d139 100644 --- a/extensions/qa-lab/src/suite.test.ts +++ b/extensions/qa-lab/src/suite.test.ts @@ -50,6 +50,40 @@ describe("qa suite", () => { expect(startLab).not.toHaveBeenCalled(); }); + it("keeps metadata-only live channel drivers on the canonical QA transport", async () => { + const create = vi.fn(); + + await expect( + qaSuiteProgressTesting.createQaSuiteTransportAdapter({ + adapterFactories: [{ id: "telegram", matches: () => true, create }], + channelDriver: "live", + outputDir: "/tmp/qa-output", + state: {} as QaLabServerHandle["state"], + transportId: "qa-channel", + }), + ).resolves.toMatchObject({ adapter: { id: "qa-channel" } }); + + expect(create).not.toHaveBeenCalled(); + }); + + it("uses a contributed live adapter when its channel is selected", async () => { + const adapter = { id: "telegram" } as QaTransportAdapter; + const create = vi.fn(async () => adapter); + + await expect( + qaSuiteProgressTesting.createQaSuiteTransportAdapter({ + adapterFactories: [{ id: "telegram", matches: () => true, create }], + channelDriver: "live", + channelId: "telegram", + outputDir: "/tmp/qa-output", + state: {} as QaLabServerHandle["state"], + transportId: "qa-channel", + }), + ).resolves.toMatchObject({ adapter }); + + expect(create).toHaveBeenCalledTimes(1); + }); + it("parses progress env booleans", () => { expect(qaSuiteProgressTesting.parseQaSuiteBooleanEnv("true")).toBe(true); expect(qaSuiteProgressTesting.parseQaSuiteBooleanEnv("on")).toBe(true); @@ -462,6 +496,11 @@ describe("qa suite", () => { it("forwards run options into isolated scenario worker params", () => { const startLab = vi.fn(); + const adapterFactory = { + id: "telegram", + matches: vi.fn(() => true), + create: vi.fn(), + }; const scenario = makeQaSuiteTestScenario("patched-control-ui", { surface: "control-ui", gatewayConfigPatch: { @@ -485,6 +524,9 @@ describe("qa suite", () => { scenario, startLab, input: { + adapterFactories: [adapterFactory], + channelId: "telegram", + adapterOptions: { repoRoot: "/repo" }, thinkingDefault: "minimal", claudeCliAuthMode: "subscription", enabledPluginIds: ["acpx"], @@ -495,6 +537,9 @@ describe("qa suite", () => { }), ).toMatchObject({ scenarioIds: ["patched-control-ui"], + adapterFactories: [adapterFactory], + channelId: "telegram", + adapterOptions: { repoRoot: "/repo" }, concurrency: 1, startLab, controlUiEnabled: true, diff --git a/extensions/qa-lab/src/suite.ts b/extensions/qa-lab/src/suite.ts index 1310151b9bf7..b81c9eab8091 100644 --- a/extensions/qa-lab/src/suite.ts +++ b/extensions/qa-lab/src/suite.ts @@ -47,6 +47,8 @@ import { createQaTransportAdapter, defaultQaSuiteConcurrencyForTransport, normalizeQaTransportId, + type QaTransportAdapterFactory, + type QaTransportFactoryContext, type QaTransportId, } from "./qa-transport-registry.js"; import type { QaTransportAdapter } from "./qa-transport.js"; @@ -115,6 +117,10 @@ type QaSuiteEnvironment = { export type QaSuiteStartLabFn = (params?: QaLabServerStartParams) => Promise; async function createQaSuiteTransportAdapter(params: { + adapterOptions?: QaSuiteRunParams["adapterOptions"]; + adapterFactories?: readonly QaTransportAdapterFactory[]; + channelDriver?: QaScorecardChannelDriver | null; + channelId?: string; channelDriverSelection?: OpenClawCrablineChannelDriverSelection | null; cleanupOnFailure?: () => Promise; outputDir: string; @@ -122,12 +128,24 @@ async function createQaSuiteTransportAdapter(params: { transportId: QaTransportId; }) { try { - return await createQaTransportAdapter({ - channelId: params.channelDriverSelection?.channel ?? params.transportId, - driver: params.channelDriverSelection ? "crabline" : params.transportId, - outputDir: params.outputDir, - state: params.state, - }); + const usesLiveAdapter = + params.channelDriver === "live" && + params.channelId !== undefined && + params.adapterFactories !== undefined; + return await createQaTransportAdapter( + { + channelId: params.channelId ?? params.channelDriverSelection?.channel ?? params.transportId, + driver: usesLiveAdapter + ? "live" + : params.channelDriverSelection + ? "crabline" + : params.transportId, + outputDir: params.outputDir, + adapterOptions: params.adapterOptions, + state: params.state, + }, + usesLiveAdapter ? params.adapterFactories : undefined, + ); } catch (error) { await params.cleanupOnFailure?.().catch(() => undefined); throw error; @@ -135,6 +153,9 @@ async function createQaSuiteTransportAdapter(params: { } export type QaSuiteRunParams = { + adapterOptions?: QaTransportFactoryContext["adapterOptions"]; + adapterFactories?: readonly QaTransportAdapterFactory[]; + channelId?: string; evidenceMode?: QaScorecardEvidenceMode; repoRoot?: string; outputDir?: string; @@ -509,6 +530,9 @@ function buildQaIsolatedScenarioWorkerParams(params: { startLab: QaSuiteStartLabFn; }): QaSuiteRunParams { return { + adapterFactories: params.input?.adapterFactories, + adapterOptions: params.input?.adapterOptions, + channelId: params.input?.channelId, repoRoot: params.repoRoot, outputDir: params.outputDir, providerMode: params.providerMode, @@ -695,6 +719,9 @@ export function buildQaSuiteSummaryJson(params: QaSuiteSummaryJsonParams): QaSui } async function runQaRuntimeParitySuite(params: { + adapterOptions?: QaSuiteRunParams["adapterOptions"]; + adapterFactories?: readonly QaTransportAdapterFactory[]; + channelId?: string; evidenceMode?: QaScorecardEvidenceMode; repoRoot: string; outputDir: string; @@ -728,7 +755,11 @@ async function runQaRuntimeParitySuite(params: { embeddedGateway: "disabled", })); const transportFactoryResult = await createQaSuiteTransportAdapter({ + adapterFactories: params.adapterFactories, + channelDriver: params.channelDriver, + channelId: params.channelId, channelDriverSelection: params.channelDriverSelection, + adapterOptions: params.adapterOptions, cleanupOnFailure: ownsLab ? () => lab.stop() : undefined, outputDir: params.outputDir, state: lab.state, @@ -781,6 +812,9 @@ async function runQaRuntimeParitySuite(params: { ); const cellStartedAt = Date.now(); const cellResult = await runQaFlowSuite({ + adapterFactories: params.adapterFactories, + channelId: params.channelId, + adapterOptions: params.adapterOptions, repoRoot: params.repoRoot, outputDir: cellOutputDir, providerMode: params.providerMode, @@ -1257,6 +1291,9 @@ export async function runQaFlowSuite(params?: QaSuiteRunParams): Promise lab.stop() : undefined, outputDir, state: lab.state, @@ -1553,7 +1594,11 @@ export async function runQaFlowSuite(params?: QaSuiteRunParams): Promise lab.stop() : undefined, outputDir, state: lab.state, @@ -1851,6 +1896,7 @@ export const qaSuiteProgressTesting = { buildQaGatewayHeapCheckpointRuntimeEnvPatch, buildQaIsolatedScenarioWorkerParams, buildQaSuiteRuntimeMetrics, + createQaSuiteTransportAdapter, formatQaSuiteRunStartProgress, buildQaRuntimeEnvPatch, mergeQaRuntimeEnvPatches, diff --git a/extensions/qa-matrix/cli.ts b/extensions/qa-matrix/cli.ts index b474c36dc046..bd0479fa4ccd 100644 --- a/extensions/qa-matrix/cli.ts +++ b/extensions/qa-matrix/cli.ts @@ -1,2 +1,2 @@ // Qa Matrix plugin module implements cli behavior. -export { qaRunnerCliRegistrations, registerMatrixQaCli } from "./src/cli.js"; +export { qaRunnerCliRegistrations } from "./src/cli.js"; diff --git a/extensions/qa-matrix/src/adapter.runtime.ts b/extensions/qa-matrix/src/adapter.runtime.ts new file mode 100644 index 000000000000..dee5676b919b --- /dev/null +++ b/extensions/qa-matrix/src/adapter.runtime.ts @@ -0,0 +1,216 @@ +// Qa Matrix plugin module implements Matrix live transport adapter behavior. +import { randomUUID } from "node:crypto"; +import path from "node:path"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import type { QaRunnerCliRegistration } from "openclaw/plugin-sdk/qa-runner-runtime"; +import { createMatrixQaClient, provisionMatrixQaRoom } from "./substrate/client.js"; +import { buildMatrixQaConfig } from "./substrate/config.js"; +import type { MatrixQaObservedEvent } from "./substrate/events.js"; +import { startMatrixQaHarness } from "./substrate/harness.runtime.js"; +import { createMatrixQaRoomObserver } from "./substrate/sync.js"; + +type AdapterFactory = NonNullable; +type FactoryContext = Parameters[0]; +type AdapterDefinition = Awaited>; + +async function waitForMatrixChannelReady( + gateway: Parameters[0]["gateway"], + accountId: string, + timeoutMs = 60_000, + pollIntervalMs = 500, +) { + const deadline = Date.now() + timeoutMs; + let lastAccounts: unknown; + while (Date.now() < deadline) { + try { + const payload = (await gateway.call( + "channels.status", + { probe: false, timeoutMs: Math.min(2_000, timeoutMs) }, + { timeoutMs: Math.min(5_000, timeoutMs) }, + )) as { + channelAccounts?: Record< + string, + Array<{ + accountId?: string; + connected?: boolean; + healthState?: string; + restartPending?: boolean; + running?: boolean; + }> + >; + }; + const accounts = payload.channelAccounts?.matrix ?? []; + lastAccounts = accounts; + const account = accounts.find((entry) => entry.accountId === accountId); + if ( + account?.running === true && + account.connected === true && + account.restartPending !== true && + account.healthState !== "degraded" + ) { + return; + } + } catch { + // Retry until the shared host readiness deadline. + } + await new Promise((resolve) => { + setTimeout(resolve, pollIntervalMs); + }); + } + throw new Error( + `matrix account "${accountId}" did not become ready; last accounts: ${JSON.stringify(lastAccounts ?? [])}`, + ); +} + +export async function createMatrixQaTransportAdapter( + context: FactoryContext, +): Promise { + const options = context.adapterOptions ?? {}; + const repoRoot = options.repoRoot?.trim() || process.cwd(); + const harness = await startMatrixQaHarness({ + outputDir: path.join(context.outputDir, "matrix-harness"), + repoRoot, + }); + const suffix = randomUUID().slice(0, 8); + let provisioning: Awaited>; + try { + provisioning = await provisionMatrixQaRoom({ + baseUrl: harness.baseUrl, + driverLocalpart: `qa-driver-${suffix}`, + observerLocalpart: `qa-observer-${suffix}`, + registrationToken: harness.registrationToken, + roomName: `OpenClaw Matrix QA ${suffix}`, + sutLocalpart: `qa-sut-${suffix}`, + }); + } catch (error) { + await harness.stop().catch(() => undefined); + throw error; + } + const accountId = options.sutAccountId?.trim() || "sut"; + const observedEvents: MatrixQaObservedEvent[] = []; + const observer = createMatrixQaRoomObserver({ + accessToken: provisioning.observer.accessToken, + baseUrl: harness.baseUrl, + observedEvents, + }); + try { + await observer.prime(); + } catch (error) { + await harness.stop().catch(() => undefined); + throw error; + } + const driverClient = createMatrixQaClient({ + accessToken: provisioning.driver.accessToken, + baseUrl: harness.baseUrl, + }); + let stopped = false; + let pollingError: Error | undefined; + let logicalConversationId = provisioning.roomId; + let logicalConversationKind: "channel" | "direct" | "group" = "channel"; + const nativeEventIds = new Map(); + const busMessageIds = new Map(); + const polling = (async () => { + for (;;) { + if (stopped) { + return; + } + const result = await observer.waitForOptionalRoomEvent({ + predicate: (event) => event.sender === provisioning.sut.userId, + roomId: provisioning.roomId, + timeoutMs: 1_000, + }); + if (!result.matched) { + continue; + } + const event = result.event; + await context.messages.addOutboundMessage({ + accountId, + to: `${logicalConversationKind}:${logicalConversationId}`, + senderId: event.sender, + text: event.body ?? "", + timestamp: event.originServerTs, + threadId: + event.relatesTo?.relType === "m.thread" && event.relatesTo.eventId + ? busMessageIds.get(event.relatesTo.eventId) + : undefined, + replyToId: event.relatesTo?.inReplyToId + ? busMessageIds.get(event.relatesTo.inReplyToId) + : undefined, + }); + } + })().catch((error: unknown) => { + if (!stopped) { + pollingError = error instanceof Error ? error : new Error(String(error)); + } + }); + + return { + id: "matrix", + label: "Matrix live", + accountId, + requiredPluginIds: ["matrix"], + supportedActions: [], + assertTransportHealthy() { + if (pollingError) { + throw pollingError; + } + }, + async sendInbound(input) { + logicalConversationId = input.conversation.id; + logicalConversationKind = input.conversation.kind; + const hasPortableMention = input.text.includes("@openclaw"); + const body = input.text.replaceAll("@openclaw", provisioning.sut.userId); + const eventId = await driverClient.sendTextMessage({ + body, + mentionUserIds: hasPortableMention ? [provisioning.sut.userId] : undefined, + replyToEventId: input.replyToId ? nativeEventIds.get(input.replyToId) : undefined, + roomId: provisioning.roomId, + threadRootEventId: input.threadId ? nativeEventIds.get(input.threadId) : undefined, + }); + const message = await context.messages.addInboundMessage({ + ...input, + accountId, + senderId: provisioning.driver.userId, + }); + nativeEventIds.set(message.id, eventId); + busMessageIds.set(eventId, message.id); + return message; + }, + resetTransport: () => { + logicalConversationId = provisioning.roomId; + logicalConversationKind = "channel"; + nativeEventIds.clear(); + busMessageIds.clear(); + }, + createGatewayConfig: () => + buildMatrixQaConfig({} as OpenClawConfig, { + driverAccessToken: provisioning.driver.accessToken, + driverUserId: provisioning.driver.userId, + homeserver: harness.baseUrl, + observerAccessToken: provisioning.observer.accessToken, + observerUserId: provisioning.observer.userId, + sutAccessToken: provisioning.sut.accessToken, + sutAccountId: accountId, + sutDeviceId: provisioning.sut.deviceId, + sutUserId: provisioning.sut.userId, + topology: provisioning.topology, + }), + waitReady: async ({ gateway, timeoutMs, pollIntervalMs }) => + await waitForMatrixChannelReady(gateway, accountId, timeoutMs, pollIntervalMs), + buildAgentDelivery: () => ({ + channel: "matrix", + to: provisioning.roomId, + replyChannel: "matrix", + replyTo: provisioning.roomId, + }), + async handleAction() { + throw new Error("Matrix live QA adapter does not implement transport actions"); + }, + createReportNotes: () => ["Uses the Matrix live adapter."], + async cleanup() { + stopped = true; + await polling.catch(() => undefined); + await harness.stop(); + }, + }; +} diff --git a/extensions/qa-matrix/src/cli.ts b/extensions/qa-matrix/src/cli.ts index f508773aa5b9..9bf4cf00a3e6 100644 --- a/extensions/qa-matrix/src/cli.ts +++ b/extensions/qa-matrix/src/cli.ts @@ -9,12 +9,16 @@ import { } from "./shared/live-transport-cli.js"; type MatrixQaCliRuntime = typeof import("./cli.runtime.js"); +type MatrixQaAdapterRuntime = typeof import("./adapter.runtime.js"); const DISABLE_MATRIX_QA_FORCE_EXIT_ENV = "OPENCLAW_QA_MATRIX_DISABLE_FORCE_EXIT"; const loadMatrixQaCliRuntime = createLazyCliRuntimeLoader( () => import("./cli.runtime.js"), ); +const loadMatrixQaAdapterRuntime = createLazyCliRuntimeLoader( + () => import("./adapter.runtime.js"), +); async function flushProcessStream(stream: NodeJS.WriteStream) { if (stream.destroyed || !stream.writable) { @@ -52,9 +56,20 @@ async function runQaMatrix(opts: LiveTransportQaCommandOptions) { } } +export const matrixQaAdapterFactory: NonNullable = + { + id: "matrix", + scenarioIds: ["channel-chat-baseline"], + matches: ({ channelId, driver }) => driver === "live" && channelId === "matrix", + async create(context) { + return await (await loadMatrixQaAdapterRuntime()).createMatrixQaTransportAdapter(context); + }, + }; + export const matrixQaCliRegistration: LiveTransportQaCliRegistration = createLiveTransportQaCliRegistration({ commandName: "matrix", + adapterFactory: matrixQaAdapterFactory, description: "Run the Docker-backed Matrix live QA lane against a disposable homeserver", outputDirHelp: "Matrix QA artifact directory", profileHelp: diff --git a/extensions/qa-matrix/src/substrate/harness.runtime.ts b/extensions/qa-matrix/src/substrate/harness.runtime.ts index a5dbd0d592db..fa46f87e1d31 100644 --- a/extensions/qa-matrix/src/substrate/harness.runtime.ts +++ b/extensions/qa-matrix/src/substrate/harness.runtime.ts @@ -40,7 +40,7 @@ type MatrixQaHarnessFiles = { registrationToken: string; }; -type MatrixQaHarness = MatrixQaHarnessFiles & { +export type MatrixQaHarness = MatrixQaHarnessFiles & { baseUrl: string; recording: MatrixQaRecordingProxy; restartService(): Promise; diff --git a/extensions/telegram/api.ts b/extensions/telegram/api.ts index c005eeeab785..9b24a3c89a4d 100644 --- a/extensions/telegram/api.ts +++ b/extensions/telegram/api.ts @@ -1,4 +1,8 @@ // Telegram API module exposes the plugin public contract. +export type { + Message as TelegramBotMessage, + Update as TelegramBotUpdate, +} from "grammy/types"; export { telegramPlugin } from "./src/channel.js"; export { telegramSetupPlugin } from "./src/channel.setup.js"; export { diff --git a/extensions/tsconfig.package-boundary.paths.json b/extensions/tsconfig.package-boundary.paths.json index 6e116c1ad332..c0a2fd799428 100644 --- a/extensions/tsconfig.package-boundary.paths.json +++ b/extensions/tsconfig.package-boundary.paths.json @@ -113,6 +113,9 @@ "@openclaw/slack/api.js": [ "../dist/plugin-sdk/extensions/slack/api.d.ts" ], + "@openclaw/telegram/api.js": [ + "../dist/plugin-sdk/extensions/telegram/api.d.ts" + ], "@openclaw/whatsapp/api.js": [ "../dist/plugin-sdk/extensions/whatsapp/api.d.ts" ], diff --git a/scripts/lib/extension-package-boundary.ts b/scripts/lib/extension-package-boundary.ts index b3d5e893fe3e..4a133525b829 100644 --- a/scripts/lib/extension-package-boundary.ts +++ b/scripts/lib/extension-package-boundary.ts @@ -84,6 +84,7 @@ export const EXTENSION_PACKAGE_BOUNDARY_BASE_PATHS = { "@openclaw/qa-channel/api.js": ["../dist/plugin-sdk/extensions/qa-channel/api.d.ts"], "@openclaw/discord/api.js": ["../dist/plugin-sdk/extensions/discord/api.d.ts"], "@openclaw/slack/api.js": ["../dist/plugin-sdk/extensions/slack/api.d.ts"], + "@openclaw/telegram/api.js": ["../dist/plugin-sdk/extensions/telegram/api.d.ts"], "@openclaw/whatsapp/api.js": ["../dist/plugin-sdk/extensions/whatsapp/api.d.ts"], "@openclaw/ai": ["../dist/plugin-sdk/packages/ai/src/index.d.ts"], "@openclaw/ai/diagnostics": ["../dist/plugin-sdk/packages/ai/src/utils/diagnostics.d.ts"], @@ -260,6 +261,7 @@ export const EXTENSION_PACKAGE_BOUNDARY_XAI_PATHS = { "openclaw/plugin-sdk/channel-secret-tts-runtime": _omitTts, "@openclaw/discord/api.js": _omitDiscord, "@openclaw/slack/api.js": _omitSlack, + "@openclaw/telegram/api.js": _omitTelegram, "@openclaw/whatsapp/api.js": _omitWhatsApp, ...rest }) => rest)(EXTENSION_PACKAGE_BOUNDARY_BASE_PATHS), diff --git a/scripts/prepare-extension-package-boundary-artifacts.mjs b/scripts/prepare-extension-package-boundary-artifacts.mjs index ef655eb3ee46..b662b860711d 100644 --- a/scripts/prepare-extension-package-boundary-artifacts.mjs +++ b/scripts/prepare-extension-package-boundary-artifacts.mjs @@ -267,6 +267,12 @@ const SLACK_DTS_INPUTS = [ ]; const SLACK_DTS_STAMP = "dist/plugin-sdk/extensions/slack/.boundary-dts.stamp"; const SLACK_DTS_REQUIRED_OUTPUTS = ["dist/plugin-sdk/extensions/slack/api.d.ts"]; +const TELEGRAM_DTS_INPUTS = [ + "extensions/telegram/api.ts", + "extensions/telegram/tsconfig.json", +]; +const TELEGRAM_DTS_STAMP = "dist/plugin-sdk/extensions/telegram/.boundary-dts.stamp"; +const TELEGRAM_DTS_REQUIRED_OUTPUTS = ["dist/plugin-sdk/extensions/telegram/api.d.ts"]; const WHATSAPP_DTS_INPUTS = [ "extensions/whatsapp/api.ts", "extensions/whatsapp/src/qa-driver.runtime.ts", @@ -748,6 +754,12 @@ async function main(argv = process.argv.slice(2)) { outputPaths: [SLACK_DTS_STAMP, ...SLACK_DTS_REQUIRED_OUTPUTS], includeFile: isRelevantTypeInput, }) && !hasMissingOutput(SLACK_DTS_REQUIRED_OUTPUTS); + const telegramDtsFresh = + isArtifactSetFresh({ + inputPaths: TELEGRAM_DTS_INPUTS, + outputPaths: [TELEGRAM_DTS_STAMP, ...TELEGRAM_DTS_REQUIRED_OUTPUTS], + includeFile: isRelevantTypeInput, + }) && !hasMissingOutput(TELEGRAM_DTS_REQUIRED_OUTPUTS); const whatsappDtsFresh = isArtifactSetFresh({ inputPaths: WHATSAPP_DTS_INPUTS, @@ -914,6 +926,37 @@ async function main(argv = process.argv.slice(2)) { } else { process.stdout.write("[whatsapp boundary dts] fresh; skipping\n"); } + if (!telegramDtsFresh) { + removeIncrementalStateForMissingOutput({ + outputPaths: TELEGRAM_DTS_REQUIRED_OUTPUTS, + tsBuildInfoPath: "dist/plugin-sdk/extensions/telegram/.tsbuildinfo", + }); + dependentSteps.push({ + label: "telegram boundary dts", + args: [ + runTsgoScript, + "-p", + "extensions/telegram/tsconfig.json", + "--declaration", + "true", + "--emitDeclarationOnly", + "true", + "--noEmit", + "false", + "--outDir", + "dist/plugin-sdk/extensions/telegram", + "--rootDir", + "extensions/telegram", + "--tsBuildInfoFile", + "dist/plugin-sdk/extensions/telegram/.tsbuildinfo", + ], + env: { OPENCLAW_TSGO_HEAVY_CHECK_LOCK_HELD: "1" }, + timeoutMs: 300_000, + stampPath: TELEGRAM_DTS_STAMP, + }); + } else { + process.stdout.write("[telegram boundary dts] fresh; skipping\n"); + } } if (prerequisiteSteps.length > 0) { diff --git a/src/plugin-sdk/qa-runner-runtime.integration.test.ts b/src/plugin-sdk/qa-runner-runtime.integration.test.ts index fc0cfa4212d8..14162054e298 100644 --- a/src/plugin-sdk/qa-runner-runtime.integration.test.ts +++ b/src/plugin-sdk/qa-runner-runtime.integration.test.ts @@ -109,6 +109,7 @@ describe("plugin-sdk qa-runner-runtime linked plugin smoke", () => { "export const qaRunnerCliRegistrations = [", " {", ' commandName: "linked",', + ' adapterFactory: { id: "linked", matches() { return true; }, async create(context) { return { id: "linked", label: "Linked", accountId: "sut", requiredPluginIds: [], supportedActions: [], async sendInbound(input) { return await context.messages.addInboundMessage(input); }, createGatewayConfig() { return {}; }, async waitReady() {}, buildAgentDelivery({ target }) { return { channel: "linked", to: target, replyChannel: "linked", replyTo: target }; }, async handleAction() {}, createReportNotes() { return []; } }; } },', " register() {}", " }", "];", @@ -132,6 +133,9 @@ describe("plugin-sdk qa-runner-runtime linked plugin smoke", () => { status: "available", registration: { commandName: "linked", + adapterFactory: expect.objectContaining({ + id: "linked", + }), register, }, }, diff --git a/src/plugin-sdk/qa-runner-runtime.test.ts b/src/plugin-sdk/qa-runner-runtime.test.ts index 9d3b4f352ca6..992910218970 100644 --- a/src/plugin-sdk/qa-runner-runtime.test.ts +++ b/src/plugin-sdk/qa-runner-runtime.test.ts @@ -130,6 +130,7 @@ describe("plugin-sdk qa-runner-runtime", () => { it("returns activated runner registrations declared in plugin manifests", async () => { const register = vi.fn((qa: Command) => qa); + const adapterFactory = { id: "matrix", matches: vi.fn(), create: vi.fn() }; loadPluginManifestRegistry.mockReturnValue({ plugins: [ { @@ -147,7 +148,7 @@ describe("plugin-sdk qa-runner-runtime", () => { diagnostics: [], }); loadBundledPluginPublicSurfaceModuleSync.mockReturnValue({ - qaRunnerCliRegistrations: [{ commandName: "matrix", register }], + qaRunnerCliRegistrations: [{ commandName: "matrix", adapterFactory, register }], }); const module = await import("./qa-runner-runtime.js"); @@ -160,6 +161,7 @@ describe("plugin-sdk qa-runner-runtime", () => { status: "available", registration: { commandName: "matrix", + adapterFactory, register, }, }, @@ -195,11 +197,41 @@ describe("plugin-sdk qa-runner-runtime", () => { ]); }); + it("keeps shipped registration-only runner contributions available", async () => { + const register = vi.fn((qa: Command) => qa); + loadPluginManifestRegistry.mockReturnValue({ + plugins: [ + { + id: "qa-legacy", + origin: "bundled", + qaRunners: [{ commandName: "legacy" }], + rootDir: "/tmp/qa-legacy", + }, + ], + diagnostics: [], + }); + loadBundledPluginPublicSurfaceModuleSync.mockReturnValue({ + qaRunnerCliRegistrations: [{ commandName: "legacy", register }], + }); + + const module = await import("./qa-runner-runtime.js"); + + expect(module.listQaRunnerCliContributions()).toEqual([ + { + pluginId: "qa-legacy", + commandName: "legacy", + status: "available", + registration: { commandName: "legacy", register }, + }, + ]); + }); + it("prefers the source bundled tree for private qa discovery in repo checkouts", async () => { const sourceRoot = makePrivateQaSourceRoot(tempDirs, "openclaw-qa-runner-root-"); resolveOpenClawPackageRootSync.mockReturnValue(sourceRoot); const register = vi.fn((qa: Command) => qa); + const adapterFactory = { id: "matrix", matches: vi.fn(), create: vi.fn() }; loadPluginManifestRegistry.mockReturnValue({ plugins: [ { @@ -212,7 +244,7 @@ describe("plugin-sdk qa-runner-runtime", () => { diagnostics: [], }); loadBundledPluginPublicSurfaceModuleSync.mockReturnValue({ - qaRunnerCliRegistrations: [{ commandName: "matrix", register }], + qaRunnerCliRegistrations: [{ commandName: "matrix", adapterFactory, register }], }); const module = await import("./qa-runner-runtime.js"); @@ -224,6 +256,7 @@ describe("plugin-sdk qa-runner-runtime", () => { status: "available", registration: { commandName: "matrix", + adapterFactory, register, }, }, @@ -284,7 +317,11 @@ describe("plugin-sdk qa-runner-runtime", () => { }); loadBundledPluginPublicSurfaceModuleSync.mockReturnValue({ qaRunnerCliRegistrations: [ - { commandName: "matrix", register: vi.fn() }, + { + commandName: "matrix", + adapterFactory: { id: "matrix", matches: vi.fn(), create: vi.fn() }, + register: vi.fn(), + }, { commandName: "extra", register: vi.fn() }, ], }); diff --git a/src/plugin-sdk/qa-runner-runtime.ts b/src/plugin-sdk/qa-runner-runtime.ts index abd7880ee5cf..a4afcedced7f 100644 --- a/src/plugin-sdk/qa-runner-runtime.ts +++ b/src/plugin-sdk/qa-runner-runtime.ts @@ -2,15 +2,111 @@ import type { Command } from "commander"; import type { PluginManifestRecord } from "../plugins/manifest-registry.js"; import { loadPluginManifestRegistry } from "../plugins/manifest-registry.js"; +import type { OpenClawConfig } from "./config-contracts.js"; import { loadBundledPluginPublicSurfaceModuleSync, tryLoadActivatedBundledPluginPublicSurfaceModuleSync, } from "./facade-runtime.js"; import { resolvePrivateQaBundledPluginsEnv } from "./private-qa-bundled-env.js"; +import type { + QaBusEditMessageInput, + QaBusInboundMessageInput, + QaBusMessage, + QaBusOutboundMessageInput, +} from "./qa-channel-protocol.js"; -/** CLI registration exported by a QA runner plugin runtime surface. */ +type QaRunnerAdapterOptions = { + repoRoot?: string; + sutAccountId?: string; + credentialSource?: string; + credentialRole?: string; +}; + +type QaRunnerMessageRecorder = { + addInboundMessage: (input: QaBusInboundMessageInput) => QaBusMessage | Promise; + addOutboundMessage: (input: QaBusOutboundMessageInput) => QaBusMessage | Promise; + editMessage: (input: QaBusEditMessageInput) => QaBusMessage | Promise; +}; + +type QaRunnerTransportAdapterDefinition = { + id: string; + label: string; + accountId: string; + requiredPluginIds: readonly string[]; + supportedActions: readonly ("delete" | "edit" | "react" | "thread-create")[]; + assertTransportHealthy?: () => void; + resetTransport?: () => void | Promise; + sendInbound: (input: QaBusInboundMessageInput) => Promise; + sendNativeCommand?: ( + input: Omit & { command: string }, + ) => Promise; + waitForOutboundSequence?: (input: { + conversationId?: string; + finalSettleMs?: number; + finalTextIncludes: string; + minimumPreviewEvents?: number; + sinceCursor?: number; + threadId?: string; + timeoutMs?: number; + }) => Promise<{ + events: Array<{ cursor: number; kind: "sent" | "edited" | "deleted"; message: QaBusMessage }>; + final: QaBusMessage; + }>; + createGatewayConfig: (params: { + baseUrl: string; + }) => Pick; + waitReady: (params: { + gateway: { + call: ( + method: string, + params?: unknown, + options?: { timeoutMs?: number }, + ) => Promise; + }; + timeoutMs?: number; + pollIntervalMs?: number; + }) => Promise; + buildAgentDelivery: (params: { target: string }) => { + channel: string; + to?: string; + replyChannel: string; + replyTo: string; + }; + createRuntimeEnvPatch?: () => NodeJS.ProcessEnv; + handleAction: (params: { + action: "delete" | "edit" | "react" | "thread-create"; + args: Record; + cfg: OpenClawConfig; + accountId?: string | null; + }) => Promise; + createReportNotes: (params: { + providerMode: "mock-openai" | "aimock" | "live-frontier"; + primaryModel: string; + alternateModel: string; + fastMode: boolean; + concurrency: number; + isolatedWorkers?: boolean; + }) => string[]; + cleanup?: () => Promise; +}; + +type QaRunnerTransportFactory = { + id: string; + scenarioIds?: readonly string[]; + matches: (context: { channelId: string; driver: string }) => boolean; + create: (context: { + adapterOptions?: QaRunnerAdapterOptions; + channelId: string; + driver: string; + messages: QaRunnerMessageRecorder; + outputDir: string; + }) => Promise; +}; + +/** CLI registration and optional transport adapter factory exported by a QA runner plugin. */ export type QaRunnerCliRegistration = { commandName: string; + adapterFactory?: QaRunnerTransportFactory; register(qa: Command): void; }; @@ -187,6 +283,17 @@ export function listQaRunnerCliContributions(): readonly QaRunnerCliContribution `QA runner plugin "${plugin.id}" declared "${runner.commandName}" in openclaw.plugin.json but did not export a matching CLI registration`, ); } + const adapterFactory = registration.adapterFactory; + if ( + adapterFactory && + (adapterFactory.id !== runner.commandName || + typeof adapterFactory.matches !== "function" || + typeof adapterFactory.create !== "function") + ) { + throw new Error( + `QA runner plugin "${plugin.id}" exported an invalid transport factory for "${runner.commandName}"`, + ); + } contributions.set(runner.commandName, { pluginId: plugin.id, commandName: runner.commandName, diff --git a/src/plugin-sdk/qa-runtime.ts b/src/plugin-sdk/qa-runtime.ts index 0305241dc2b6..2af200f968af 100644 --- a/src/plugin-sdk/qa-runtime.ts +++ b/src/plugin-sdk/qa-runtime.ts @@ -9,6 +9,7 @@ import { formatErrorMessage } from "./error-runtime.js"; import { loadBundledPluginPublicSurfaceModuleSync } from "./facade-runtime.js"; import { resolvePrivateQaBundledPluginsEnv } from "./private-qa-bundled-env.js"; import { runExec } from "./process-runtime.js"; +import type { QaRunnerCliRegistration } from "./qa-runner-runtime.js"; import { fetchWithSsrFGuard } from "./ssrf-runtime.js"; import { normalizeStringEntries } from "./string-coerce-runtime.js"; @@ -90,10 +91,7 @@ type LiveTransportQaCommanderOptions = { }; /** Commander registration hook for one live-transport QA subcommand. */ -export type LiveTransportQaCliRegistration = { - commandName: string; - register(qa: Command): void; -}; +export type LiveTransportQaCliRegistration = QaRunnerCliRegistration; /** Help text customizations for live credential source and role flags. */ export type LiveTransportQaCredentialCliOptions = { @@ -115,6 +113,7 @@ export type LiveTransportQaCliRegistrationOptions = { allowFailuresHelp?: string; scenarioHelp: string; sutAccountHelp: string; + adapterFactory?: QaRunnerCliRegistration["adapterFactory"]; run: (opts: LiveTransportQaCommandOptions) => Promise; }; @@ -156,6 +155,7 @@ function mapLiveTransportQaCommanderOptions( function registerLiveTransportQaCli( params: LiveTransportQaCliRegistrationOptions & { qa: Command; + run: (opts: LiveTransportQaCommandOptions) => Promise; }, ) { const command = params.qa @@ -209,6 +209,7 @@ export function createLiveTransportQaCliRegistration( ): LiveTransportQaCliRegistration { return { commandName: params.commandName, + adapterFactory: params.adapterFactory, register(qa: Command) { registerLiveTransportQaCli({ ...params,