From 33ea7ffa547f0fe512301f5a13118348fc10ce7a Mon Sep 17 00:00:00 2001 From: Dallin Romney Date: Sun, 9 Aug 2026 06:11:03 +0800 Subject: [PATCH] fix(ci): isolate package Telegram QA harness (#120193) * fix(qa): isolate package Telegram harness Keep private QA source, dependencies, taxonomy, and SDK dist in the trusted harness while the installed candidate owns its CLI, Gateway runtime, and persisted mock auth. Preserve the documented package RTT canary after taxonomy selection. Co-authored-by: Dallin Romney Co-authored-by: Vincent Koc * fix(qa): export private QA harness SDK entries Canonicalize the QA-only plugin SDK entries shared by the private build and package Telegram harness manifest so qa-runtime and qa-lab resolve from trusted dist. * fix(qa): expose private runtime to package harness * fix(qa): surface Telegram observer conflicts * fix(qa): accept separate preview and final messages * test(qa): exercise Telegram poll delay contract --------- Co-authored-by: Peter Steinberger Co-authored-by: Vincent Koc --- docs/help/testing.md | 13 +- extensions/qa-lab/src/gateway-child.test.ts | 164 ++++++++++- extensions/qa-lab/src/gateway-child.ts | 82 +++++- .../shared/scenario-selection.test.ts | 14 + .../telegram/adapter.runtime.test.ts | 206 ++++++++++++++ .../telegram/adapter.runtime.ts | 74 ++++- .../live-transports/telegram/profiles.test.ts | 9 + .../telegram/telegram-api.runtime.test.ts | 151 +++++++++- .../telegram/telegram-api.runtime.ts | 104 ++++++- .../qa-lab/src/providers/shared/mock-auth.ts | 2 +- extensions/qa-lab/src/qa-transport.test.ts | 89 ++++++ extensions/qa-lab/src/qa-transport.ts | 259 +++++++++++------- .../src/scenario-catalog-channels.test.ts | 7 + qa/scenarios/channels/channel-canary.yaml | 3 + .../channels/channel-message-flows.yaml | 1 + .../lib/npm-telegram-live/prepare-package.mjs | 34 ++- scripts/e2e/npm-telegram-live-docker.sh | 125 ++++----- scripts/e2e/npm-telegram-live-runner.ts | 42 ++- src/plugin-sdk/qa-runner-runtime.ts | 1 + test/scripts/lint-suppressions.test.ts | 1 + test/scripts/npm-telegram-live.test.ts | 195 +++++++------ 21 files changed, 1276 insertions(+), 300 deletions(-) diff --git a/docs/help/testing.md b/docs/help/testing.md index e75f05c97efa..19a16a9d16cb 100644 --- a/docs/help/testing.md +++ b/docs/help/testing.md @@ -266,10 +266,10 @@ inside every shard. onboarding, configures Telegram through the installed CLI, then reuses the live Telegram QA lane with that installed package as the SUT Gateway. - - The wrapper mounts only the `qa-lab` harness source from the checkout; - the installed package owns `dist`, `openclaw/plugin-sdk`, and bundled - plugin runtime, so the lane does not mix current checkout plugins into - the package under test. + - The trusted checkout owns the QA harness source, taxonomy, scenarios, + dependencies, and private SDK build. The installed package remains the + absolute CLI, Gateway, and bundled-plugin runtime under test, and its CLI + writes the package candidate's persisted auth state. - Defaults to `OPENCLAW_NPM_TELEGRAM_PACKAGE_SPEC=openclaw@beta`; set `OPENCLAW_NPM_TELEGRAM_PACKAGE_TGZ=/path/to/openclaw-current.tgz` or `OPENCLAW_CURRENT_PACKAGE_TGZ` to test a resolved local tarball instead @@ -280,7 +280,10 @@ inside every shard. `OPENCLAW_NPM_TELEGRAM_RTT_TIMEOUT_MS`, or `OPENCLAW_NPM_TELEGRAM_RTT_MAX_FAILURES` to tune the run. `OPENCLAW_NPM_TELEGRAM_RTT_CHECKS` selects the Telegram QA scenario to - sample; the supported RTT target is `channel-canary`. + sample; the supported RTT target is `channel-canary`. The package runner + promotes that portable canary once to the first position, making + canary+RTT the preflight before the remaining taxonomy-backed fail-fast + release scenarios. - Uses the same Telegram env credentials or Convex credential source as `pnpm openclaw qa telegram`. For CI/release automation, set `OPENCLAW_NPM_TELEGRAM_CREDENTIAL_SOURCE=convex` plus diff --git a/extensions/qa-lab/src/gateway-child.test.ts b/extensions/qa-lab/src/gateway-child.test.ts index 4ecb8035787a..3d26c45cf892 100644 --- a/extensions/qa-lab/src/gateway-child.test.ts +++ b/extensions/qa-lab/src/gateway-child.test.ts @@ -126,6 +126,65 @@ async function writeTempProviderConfig(value: unknown) { return configPath; } +async function writePackagedGatewayFixture(root: string): Promise { + const fixturePath = path.join(root, "packaged-gateway-fixture.mjs"); + await writeFile( + fixturePath, + `import fs from "node:fs"; +import path from "node:path"; + +const args = process.argv.slice(2); +const recordPath = process.env.QA_RECORD_PATH; +const configPath = process.env.OPENCLAW_CONFIG_PATH; +const stateDir = process.env.OPENCLAW_STATE_DIR; +if (!recordPath || !configPath || !stateDir) { + throw new Error("missing fixture environment"); +} +const record = (value) => fs.appendFileSync(recordPath, JSON.stringify(value) + "\\n"); +if (args[0] === "models") { + let stdin = ""; + process.stdin.setEncoding("utf8"); + for await (const chunk of process.stdin) stdin += chunk; + const provider = args[args.indexOf("--provider") + 1]; + record({ + kind: "auth", + args, + stdin, + dbExists: fs.existsSync(path.join(stateDir, "agents", "qa", "agent", "openclaw-agent.sqlite")), + env: { + OPENCLAW_CLI: process.env.OPENCLAW_CLI, + OPENCLAW_CONFIG_PATH: configPath, + OPENCLAW_STATE_DIR: stateDir, + }, + }); + if (process.env.QA_FAIL_PROVIDER === provider) { + process.stderr.write("Authorization: Bearer " + stdin.trim()); + process.exit(9); + } + const config = JSON.parse(fs.readFileSync(configPath, "utf8")); + config.fixtureProfiles = [...(config.fixtureProfiles ?? []), provider]; + fs.writeFileSync(configPath, JSON.stringify(config)); + process.exit(0); +} +const config = JSON.parse(fs.readFileSync(configPath, "utf8")); +record({ kind: "gateway", args, fixtureProfiles: config.fixtureProfiles }); +process.stderr.write("fixture gateway exit"); +process.exit(17); +`, + "utf8", + ); + return fixturePath; +} + +async function readJsonLines(filePath: string): Promise>> { + const contents = await readFile(filePath, "utf8"); + return contents + .trim() + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line) as Record); +} + describe("runQaGatewayCliCommand", () => { it("runs CLI commands with the Gateway fixture environment", async () => { const output = await testing.runQaGatewayCliCommand({ @@ -522,7 +581,7 @@ describe("buildQaRuntimeEnv", () => { }, transportBaseUrl: "http://127.0.0.1:43123", }), - ).rejects.toThrow(/gateway failed to spawn: .*ENOENT/u); + ).rejects.toThrow(/installed package mock auth bootstrap failed for openai: .*ENOENT/u); await expect(readdir(preferredTempParent)).resolves.toStrictEqual([]); await expect(readdir(commandTempParent)).resolves.toStrictEqual([]); @@ -1365,6 +1424,109 @@ describe("buildQaRuntimeEnv", () => { } }); + it("lets an explicit packaged command own mock auth state before gateway spawn", async () => { + const fixtureRoot = await tempDirs.makeTempDir("qa-packaged-auth-"); + const tempParentDir = path.join(fixtureRoot, "gateway-temp"); + const recordPath = path.join(fixtureRoot, "commands.jsonl"); + const fixturePath = await writePackagedGatewayFixture(fixtureRoot); + await mkdir(tempParentDir); + + await expect( + startQaGatewayChild({ + repoRoot: process.cwd(), + command: { + executablePath: process.execPath, + argsPrefix: [fixturePath], + tempParentDir, + usePackagedPlugins: true, + }, + providerMode: "mock-openai", + transportBaseUrl: "http://127.0.0.1:43123", + runtimeEnvPatch: { QA_RECORD_PATH: recordPath }, + }), + ).rejects.toThrow("fixture gateway exit"); + + const records = await readJsonLines(recordPath); + const authRecords = records.filter((record) => record.kind === "auth"); + expect(authRecords).toHaveLength(2); + expect(authRecords.map((record) => record.args)).toEqual([ + [ + "models", + "auth", + "--agent", + "qa", + "paste-api-key", + "--provider", + "openai", + "--profile-id", + "qa-mock-openai", + ], + [ + "models", + "auth", + "--agent", + "qa", + "paste-api-key", + "--provider", + "anthropic", + "--profile-id", + "qa-mock-anthropic", + ], + ]); + for (const record of authRecords) { + expect(record.dbExists).toBe(false); + expect(record.stdin).toMatch(/^sk-qa-mock-[a-f0-9]{32}\n$/u); + expect(record.env).toMatchObject({ + OPENCLAW_CLI: "1", + }); + } + expect(records.at(-1)).toMatchObject({ + kind: "gateway", + fixtureProfiles: ["openai", "anthropic"], + }); + }); + + it("blocks packaged gateway spawn when candidate auth bootstrap fails", async () => { + const fixtureRoot = await tempDirs.makeTempDir("qa-packaged-auth-fail-"); + const tempParentDir = path.join(fixtureRoot, "gateway-temp"); + const recordPath = path.join(fixtureRoot, "commands.jsonl"); + const fixturePath = await writePackagedGatewayFixture(fixtureRoot); + await mkdir(tempParentDir); + + const result = startQaGatewayChild({ + repoRoot: process.cwd(), + command: { + executablePath: process.execPath, + argsPrefix: [fixturePath], + tempParentDir, + usePackagedPlugins: true, + }, + providerMode: "mock-openai", + transportBaseUrl: "http://127.0.0.1:43123", + runtimeEnvPatch: { + QA_FAIL_PROVIDER: "openai", + QA_RECORD_PATH: recordPath, + }, + }); + + const error = await result.catch((caught: unknown) => caught); + expect(error).toBeInstanceOf(Error); + if (!(error instanceof Error)) { + throw new Error("expected package auth bootstrap error"); + } + expect(error.message).toContain( + "installed package mock auth bootstrap failed for openai: OpenClaw CLI exited 9: Authorization: Bearer ", + ); + const records = await readJsonLines(recordPath); + expect(records).toHaveLength(1); + expect(records[0]).toMatchObject({ kind: "auth", dbExists: false }); + const submittedKey = String(records[0]?.stdin).trim(); + expect(submittedKey).toMatch(/^sk-qa-mock-[a-f0-9]{32}$/u); + expect(error.message).not.toContain(submittedKey); + expect(String(error.cause)).not.toContain(submittedKey); + expect(records.some((record) => record.kind === "gateway")).toBe(false); + }); + it("stages mock profiles only for the requested agents and providers when callers override the defaults", async () => { const stateDir = await tempDirs.makeTempDir("qa-mock-auth-override-"); diff --git a/extensions/qa-lab/src/gateway-child.ts b/extensions/qa-lab/src/gateway-child.ts index 71707286d972..13491f2f4fb3 100644 --- a/extensions/qa-lab/src/gateway-child.ts +++ b/extensions/qa-lab/src/gateway-child.ts @@ -69,7 +69,7 @@ import { stageQaLiveApiKeyProfiles, stageQaLiveAnthropicSetupToken, } from "./providers/live-frontier/auth.js"; -import { stageQaMockAuthProfiles } from "./providers/shared/mock-auth.js"; +import { buildQaMockProfileId, stageQaMockAuthProfiles } from "./providers/shared/mock-auth.js"; import { listMockCodexModelInfos } from "./providers/shared/mock-model-config.js"; import { seedQaAgentWorkspace } from "./qa-agent-workspace.js"; import { buildQaGatewayConfig, type QaThinkingLevel } from "./qa-gateway-config.js"; @@ -88,6 +88,7 @@ const QA_GATEWAY_CHILD_GRACEFUL_SHUTDOWN_TIMEOUT_MS = 30_000; // Loaded Docker runners can take several seconds to reap a force-killed process group. const QA_GATEWAY_CHILD_FORCE_SHUTDOWN_TIMEOUT_MS = 10_000; const QA_GATEWAY_LOG_CLOSE_TIMEOUT_MS = 5_000; +const QA_PACKAGE_AUTH_FAILURE_MAX_CHARS = 2_048; const QA_MOCK_OPENAI_API_KEY = ["qa", "mock", "openai", "key"].join("-"); const QA_GATEWAY_CHILD_BLOCKED_SECRET_ENV_VARS = Object.freeze([ "OPENCLAW_QA_CONVEX_SECRET_CI", @@ -179,13 +180,64 @@ async function runQaGatewayCliCommand(params: { args: readonly string[]; cwd: string; env: NodeJS.ProcessEnv; + stdin?: string; }): Promise { + const hasStdin = params.stdin !== undefined; const child = spawn(params.executablePath, [...params.argsPrefix, ...params.args], { cwd: params.cwd, env: { ...params.env, OPENCLAW_CLI: "1" }, - stdio: ["ignore", "pipe", "pipe"], + stdio: [hasStdin ? "pipe" : "ignore", "pipe", "pipe"], }); - return await readQaGatewayCliCommand(child); + const result = readQaGatewayCliCommand(child); + if (hasStdin) { + child.stdin?.once("error", () => {}); + child.stdin?.end(params.stdin); + } + return await result; +} + +function createQaPackagedMockApiKey(): string { + const prefix = ["s", "k"].join(""); + return `${prefix}-${["qa", "mock", randomUUID().replaceAll("-", "")].join("-")}`; +} + +async function stageQaPackagedMockAuthProfiles(params: { + command: QaGatewayChildCommand; + cwd: string; + env: NodeJS.ProcessEnv; + providers: readonly string[]; +}): Promise { + for (const provider of uniqueStrings(params.providers)) { + try { + await runQaGatewayCliCommand({ + executablePath: params.command.executablePath, + argsPrefix: params.command.argsPrefix ?? [], + args: [ + "models", + "auth", + "--agent", + "qa", + "paste-api-key", + "--provider", + provider, + "--profile-id", + buildQaMockProfileId(provider), + ], + cwd: params.command.cwd ?? params.cwd, + env: params.env, + stdin: `${createQaPackagedMockApiKey()}\n`, + }); + } catch (error) { + const errorMessage = toErrorObject(error, "installed package auth command failed").message; + const details = sliceUtf16Safe( + redactQaGatewayDebugText(errorMessage), + 0, + QA_PACKAGE_AUTH_FAILURE_MAX_CHARS, + ); + // oxlint-disable-next-line preserve-caught-error -- Candidate CLI errors can contain the submitted API key; only the redacted message crosses this boundary. + throw new Error(`installed package mock auth bootstrap failed for ${provider}: ${details}`); + } + } } type QaChildFailure = { @@ -1189,6 +1241,7 @@ export async function startQaGatewayChild(params: { const gatewayCommand = params.command ?? (params.useRepoCli ? resolveQaGatewayChildCommand(params.repoRoot) : undefined); + const usesPackagedCandidate = params.command?.usePackagedPlugins === true; const gatewayExecutablePath = gatewayCommand?.executablePath; const gatewayArgsPrefix = gatewayCommand?.argsPrefix ?? []; const gatewayArgsSuffix = gatewayCommand?.argsSuffix ?? []; @@ -1282,12 +1335,14 @@ export async function startQaGatewayChild(params: { }); const mockAuthProviders = getQaProvider(providerMode).mockAuthProviders; if (mockAuthProviders && mockAuthProviders.length > 0) { - cfg = await stageQaMockAuthProfiles({ - cfg, - stateDir, - agentIds: params.mockAuthAgentIds, - providers: mockAuthProviders, - }); + if (!usesPackagedCandidate) { + cfg = await stageQaMockAuthProfiles({ + cfg, + stateDir, + agentIds: params.mockAuthAgentIds, + providers: mockAuthProviders, + }); + } } return params.mutateConfig ? params.mutateConfig(cfg) : cfg; }; @@ -1506,6 +1561,15 @@ export async function startQaGatewayChild(params: { encoding: "utf8", mode: 0o600, }); + const mockAuthProviders = getQaProvider(providerMode).mockAuthProviders; + if (usesPackagedCandidate && gatewayCommand && mockAuthProviders?.length) { + await stageQaPackagedMockAuthProfiles({ + command: gatewayCommand, + cwd: gatewayCwd, + env, + providers: mockAuthProviders, + }); + } } if (!env) { throw new Error("qa gateway runtime env not initialized"); diff --git a/extensions/qa-lab/src/live-transports/shared/scenario-selection.test.ts b/extensions/qa-lab/src/live-transports/shared/scenario-selection.test.ts index 249f5e0542d6..85ae5e08c5e5 100644 --- a/extensions/qa-lab/src/live-transports/shared/scenario-selection.test.ts +++ b/extensions/qa-lab/src/live-transports/shared/scenario-selection.test.ts @@ -138,6 +138,7 @@ describe("live transport QA scenario selection", () => { it.each([ { channelId: "matrix", scenarioId: "thread-follow-up" }, + { channelId: "telegram", scenarioId: "channel-canary" }, { channelId: "telegram", scenarioId: "channel-message-flows" }, ] as const)( "keeps $scenarioId eligible through both $channelId drivers", @@ -154,4 +155,17 @@ describe("live transport QA scenario selection", () => { expect(selectForDriver("crabline")).toEqual([scenarioId]); }, ); + + it("rejects the shared channel canary on unsupported Discord drivers", () => { + expect(() => + resolveCatalogLiveTransportQaScenarioIds({ + ...MOCK_LANE, + channelId: "discord", + channelDriver: "crabline", + scenarioIds: ["channel-canary"], + }), + ).toThrow( + "selected QA scenario(s) do not match the current QA lane: channel-canary (channel=qa-channel|telegram)", + ); + }); }); diff --git a/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.test.ts b/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.test.ts index a1f224f938ef..f8ed45ce68ac 100644 --- a/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.test.ts +++ b/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.test.ts @@ -10,6 +10,7 @@ const mocks = vi.hoisted(() => ({ leaseHeartbeat: vi.fn(), leaseRelease: vi.fn(), shouldRetainQaGatewayCredentialLease: vi.fn(), + waitForTelegramPollRetryDelay: vi.fn(), waitForTelegramChannelRunning: vi.fn(), })); @@ -30,10 +31,12 @@ vi.mock("./telegram-api.runtime.js", async (importOriginal) => ({ ...(await importOriginal()), callTelegramApi: mocks.callTelegramApi, flushTelegramUpdates: mocks.flushTelegramUpdates, + waitForTelegramPollRetryDelay: mocks.waitForTelegramPollRetryDelay, waitForTelegramChannelRunning: mocks.waitForTelegramChannelRunning, })); import { createTelegramQaTransportAdapter } from "./adapter.runtime.js"; +import { TelegramQaApiError } from "./telegram-api.runtime.js"; describe("Telegram QA transport adapter", () => { beforeEach(() => { @@ -50,6 +53,7 @@ describe("Telegram QA transport adapter", () => { }); mocks.flushTelegramUpdates.mockResolvedValue(0); mocks.shouldRetainQaGatewayCredentialLease.mockResolvedValue(false); + mocks.waitForTelegramPollRetryDelay.mockResolvedValue(undefined); }); it("rejects credentials that do not identify two distinct bots", async () => { @@ -68,6 +72,208 @@ describe("Telegram QA transport adapter", () => { expect(mocks.flushTelegramUpdates).not.toHaveBeenCalled(); }); + it("surfaces a duplicate getUpdates conflict without retrying it", async () => { + let getMeCalls = 0; + mocks.callTelegramApi.mockImplementation(async (_token: string, method: string) => { + if (method === "getMe") { + getMeCalls += 1; + return getMeCalls === 1 + ? { id: 1, is_bot: true, first_name: "driver", username: "driver_bot" } + : { id: 2, is_bot: true, first_name: "sut", username: "sut_bot" }; + } + throw new TelegramQaApiError( + "getUpdates", + 409, + "Conflict: terminated by other getUpdates request; make sure that only one bot instance is running", + undefined, + 409, + ); + }); + + const adapter = await createTelegramQaTransportAdapter({ + adapterOptions: {}, + messages: {}, + } as never); + + await vi.waitFor(() => expect(() => adapter.assertTransportHealthy?.()).toThrow(/Conflict/u)); + expect( + mocks.callTelegramApi.mock.calls.filter(([, method]) => method === "getUpdates"), + ).toHaveLength(1); + const diagnostics = adapter.describeTransportState?.() ?? ""; + expect(diagnostics).toContain("polls=1"); + expect(diagnostics).toContain( + "terminal error={name=TelegramQaApiError,method=getUpdates,error_code=409,status=409}", + ); + expect(diagnostics).not.toContain("placeholder"); + expect(diagnostics).not.toContain("-100123"); + + await adapter.cleanup?.(); + await adapter.cleanupAfterGatewayStop?.(); + }); + + it("backs off consecutive observer failures and resets after a successful poll", async () => { + let resolveTerminalPoll: ((updates: unknown[]) => void) | undefined; + const terminalPoll = new Promise((resolve) => { + resolveTerminalPoll = resolve; + }); + const failures = [ + new TelegramQaApiError("getUpdates", 502, "Bad Gateway", undefined, 502), + new Error("fetch failed"), + new TelegramQaApiError("getUpdates", 500, "Server Error", undefined, 500), + ] as const; + let getMeCalls = 0; + let pollCalls = 0; + mocks.callTelegramApi.mockImplementation(async (_token: string, method: string) => { + if (method === "getMe") { + getMeCalls += 1; + return getMeCalls === 1 + ? { id: 1, is_bot: true, first_name: "driver", username: "driver_bot" } + : { id: 2, is_bot: true, first_name: "sut", username: "sut_bot" }; + } + pollCalls += 1; + if (pollCalls === 1) { + throw failures[0]; + } + if (pollCalls === 2) { + throw failures[1]; + } + if (pollCalls === 3) { + return []; + } + if (pollCalls === 4) { + throw failures[2]; + } + return await terminalPoll; + }); + + const adapter = await createTelegramQaTransportAdapter({ + adapterOptions: {}, + messages: {}, + } as never); + + await vi.waitFor(() => expect(mocks.waitForTelegramPollRetryDelay).toHaveBeenCalledTimes(3)); + expect(mocks.waitForTelegramPollRetryDelay.mock.calls).toEqual([ + [failures[0], 1, expect.any(AbortSignal)], + [failures[1], 2, expect.any(AbortSignal)], + [failures[2], 1, expect.any(AbortSignal)], + ]); + + const cleanup = adapter.cleanup?.(); + resolveTerminalPoll?.([]); + await cleanup; + await adapter.cleanupAfterGatewayStop?.(); + }); + + it("aborts an observer retry delay during cleanup", async () => { + let getMeCalls = 0; + mocks.callTelegramApi.mockImplementation(async (_token: string, method: string) => { + if (method === "getMe") { + getMeCalls += 1; + return getMeCalls === 1 + ? { id: 1, is_bot: true, first_name: "driver", username: "driver_bot" } + : { id: 2, is_bot: true, first_name: "sut", username: "sut_bot" }; + } + throw new Error("fetch failed"); + }); + mocks.waitForTelegramPollRetryDelay.mockImplementation( + async (_error: unknown, _attempt: number, signal: AbortSignal) => + await new Promise((_resolve, reject) => { + signal.addEventListener("abort", () => reject(new Error("aborted")), { once: true }); + }), + ); + const adapter = await createTelegramQaTransportAdapter({ + adapterOptions: {}, + messages: {}, + } as never); + await vi.waitFor(() => expect(mocks.waitForTelegramPollRetryDelay).toHaveBeenCalledOnce()); + const signal = mocks.waitForTelegramPollRetryDelay.mock.calls[0]?.[2] as AbortSignal; + + await adapter.cleanup?.(); + + expect(signal.aborted).toBe(true); + expect(() => adapter.assertTransportHealthy?.()).not.toThrow(); + await adapter.cleanupAfterGatewayStop?.(); + }); + + it("summarizes matched and filtered update kinds without native identifiers", async () => { + const pollResolvers: Array<(updates: unknown[]) => void> = []; + let getMeCalls = 0; + mocks.callTelegramApi.mockImplementation(async (_token: string, method: string) => { + if (method === "getMe") { + getMeCalls += 1; + return getMeCalls === 1 + ? { id: 1, is_bot: true, first_name: "driver", username: "driver_bot" } + : { id: 2, is_bot: true, first_name: "sut", username: "sut_bot" }; + } + if (method === "getUpdates") { + return await new Promise((resolve) => { + pollResolvers.push(resolve); + }); + } + throw new Error(`unexpected Telegram API method: ${method}`); + }); + const addOutboundMessage = vi.fn().mockResolvedValue({ id: "out-1" }); + const adapter = await createTelegramQaTransportAdapter({ + adapterOptions: {}, + messages: { + addOutboundMessage, + editMessage: vi.fn(), + }, + } as never); + + await vi.waitFor(() => expect(pollResolvers).toHaveLength(1)); + pollResolvers[0]?.([ + { update_id: 101 }, + { + update_id: 102, + message: { + message_id: 201, + date: 100, + chat: { id: -100999 }, + from: { id: 2, is_bot: true }, + text: "wrong chat", + }, + }, + { + update_id: 103, + message: { + message_id: 202, + date: 100, + chat: { id: -100123 }, + from: { id: 3, is_bot: true }, + text: "wrong sender", + }, + }, + { + update_id: 104, + edited_message: { + message_id: 203, + date: 100, + chat: { id: -100123 }, + from: { id: 2, is_bot: true }, + text: "private matched content", + }, + }, + ]); + await vi.waitFor(() => expect(addOutboundMessage).toHaveBeenCalledOnce()); + await vi.waitFor(() => expect(pollResolvers).toHaveLength(2)); + + const diagnostics = adapter.describeTransportState?.() ?? ""; + expect(diagnostics).toContain("polls=2"); + expect(diagnostics).toContain("updates=4"); + expect(diagnostics).toContain("filtered=3"); + expect(diagnostics).toContain("matched=1"); + expect(diagnostics).toContain("update kinds=[other,message,edited_message]"); + expect(diagnostics).not.toMatch( + /-100123|10[1-4]|20[1-3]|wrong chat|wrong sender|private matched content/u, + ); + + const cleanup = adapter.cleanup?.(); + pollResolvers[1]?.([]); + await cleanup; + await adapter.cleanupAfterGatewayStop?.(); + }); + it("maps native sends, replies, edits, and cleanup inside the adapter", async () => { const pollResolvers: Array<(updates: unknown[]) => void> = []; let getMeCalls = 0; diff --git a/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts b/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts index 60bd468c4da5..7eb09d6d89b7 100644 --- a/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts +++ b/extensions/qa-lab/src/live-transports/telegram/adapter.runtime.ts @@ -17,6 +17,7 @@ import { normalizeTelegramObservedMessage, parseTelegramQaCredentialPayload, resolveTelegramQaRuntimeEnv, + TelegramQaApiError, waitForTelegramChannelRunning, waitForTelegramPollRetryDelay, type TelegramBotIdentity, @@ -28,6 +29,45 @@ type AdapterFactory = NonNullable; type FactoryContext = Parameters[0]; type AdapterDefinition = Awaited>; +const TELEGRAM_QA_DIAGNOSTIC_COUNT_LIMIT = 9_999; +type TelegramQaObserverState = { + filteredCount: number; + matchedCount: number; + pollCount: number; + relevantUpdateKinds: Set<"edited_message" | "message" | "other">; + terminalError?: Error; + updateCount: number; +}; + +function renderTelegramQaDiagnosticCount(value: number) { + return value > TELEGRAM_QA_DIAGNOSTIC_COUNT_LIMIT + ? `${TELEGRAM_QA_DIAGNOSTIC_COUNT_LIMIT}+` + : String(value); +} + +function describeTelegramQaTerminalError(error: Error | undefined) { + if (!error) { + return "none"; + } + if (error instanceof TelegramQaApiError) { + return `{name=TelegramQaApiError,method=${error.method},error_code=${error.error_code},status=${error.status}}`; + } + return `{name=${error.name === "Error" ? "Error" : "unknown"}}`; +} + +function describeTelegramQaObserverState(state: TelegramQaObserverState) { + const updateKinds = + state.relevantUpdateKinds.size > 0 ? [...state.relevantUpdateKinds] : ["none"]; + return [ + `telegram observer polls=${renderTelegramQaDiagnosticCount(state.pollCount)}`, + `updates=${renderTelegramQaDiagnosticCount(state.updateCount)}`, + `filtered=${renderTelegramQaDiagnosticCount(state.filteredCount)}`, + `matched=${renderTelegramQaDiagnosticCount(state.matchedCount)}`, + `update kinds=[${updateKinds.join(",")}]`, + `terminal error=${describeTelegramQaTerminalError(state.terminalError)}`, + ].join("; "); +} + function renderTelegramQaInboundText( input: { text: string; nativeCommand?: { name: string } }, botUsername: string, @@ -98,18 +138,27 @@ export async function createTelegramQaTransportAdapter( } const accountId = options.sutAccountId?.trim() || "sut"; let stopped = false; - let pollingError: Error | undefined; + const observerState: TelegramQaObserverState = { + filteredCount: 0, + matchedCount: 0, + pollCount: 0, + relevantUpdateKinds: new Set(), + updateCount: 0, + }; let logicalConversationId = runtimeEnv.groupId; let logicalConversationKind: "channel" | "direct" | "group" = "channel"; const nativeMessageIds = new Map(); const busMessageIds = new Map(); + const pollingAbort = new AbortController(); const poll = async () => { + let retryAttempt = 0; for (;;) { if (stopped) { return; } let updates: TelegramUpdate[]; try { + observerState.pollCount += 1; updates = await callTelegramApi( runtimeEnv.driverToken, "getUpdates", @@ -120,10 +169,16 @@ export async function createTelegramQaTransportAdapter( if (!isRecoverableTelegramQaPollError(error)) { throw error; } - await waitForTelegramPollRetryDelay(); + retryAttempt += 1; + await waitForTelegramPollRetryDelay(error, retryAttempt, pollingAbort.signal); continue; } + retryAttempt = 0; + observerState.updateCount += updates.length; for (const update of updates) { + observerState.relevantUpdateKinds.add( + update.edited_message ? "edited_message" : update.message ? "message" : "other", + ); offset = Math.max(offset, update.update_id + 1); const message = normalizeTelegramObservedMessage(update); if ( @@ -131,8 +186,10 @@ export async function createTelegramQaTransportAdapter( message.chatId !== Number(runtimeEnv.groupId) || message.senderId !== sutIdentity.id ) { + observerState.filteredCount += 1; continue; } + observerState.matchedCount += 1; const existingMessageId = busMessageIds.get(message.messageId); if (update.edited_message && existingMessageId) { await context.messages.editMessage({ @@ -163,7 +220,7 @@ export async function createTelegramQaTransportAdapter( }; const polling = poll().catch((error: unknown) => { if (!stopped) { - pollingError = error instanceof Error ? error : new Error(String(error)); + observerState.terminalError = error instanceof Error ? error : new Error("unknown error"); } }); return { @@ -173,11 +230,12 @@ export async function createTelegramQaTransportAdapter( requiredPluginIds: ["telegram"], supportedActions: [], assertTransportHealthy() { - if (pollingError) { - throw pollingError; + if (observerState.terminalError) { + throw observerState.terminalError; } heartbeat.throwIfFailed(); }, + describeTransportState: () => describeTelegramQaObserverState(observerState), async sendInbound(input) { heartbeat.throwIfFailed(); logicalConversationId = input.conversation.id; @@ -216,6 +274,11 @@ export async function createTelegramQaTransportAdapter( logicalConversationKind = "channel"; nativeMessageIds.clear(); busMessageIds.clear(); + observerState.pollCount = 0; + observerState.updateCount = 0; + observerState.filteredCount = 0; + observerState.matchedCount = 0; + observerState.relevantUpdateKinds.clear(); }, createGatewayConfig: () => buildTelegramQaConfig({} as OpenClawConfig, { @@ -243,6 +306,7 @@ export async function createTelegramQaTransportAdapter( createReportNotes: () => ["Runs through the Telegram live adapter and shared QA suite host."], async cleanup() { stopped = true; + pollingAbort.abort(new Error("Telegram QA observer stopped")); await polling.catch(() => undefined); }, async cleanupAfterGatewayStop() { diff --git a/extensions/qa-lab/src/live-transports/telegram/profiles.test.ts b/extensions/qa-lab/src/live-transports/telegram/profiles.test.ts index 22c7cc6b6993..700b3a6c07ff 100644 --- a/extensions/qa-lab/src/live-transports/telegram/profiles.test.ts +++ b/extensions/qa-lab/src/live-transports/telegram/profiles.test.ts @@ -8,6 +8,7 @@ describe("Telegram QA profiles", () => { (providerMode) => { const scenarioIds = resolveTelegramQaScenarioIds({ providerMode }); + expect(scenarioIds).toContain("channel-canary"); expect(scenarioIds).toContain("telegram-other-bot-command-gating"); expect(scenarioIds).not.toContain("telegram-startup-getme-live"); expect(() => @@ -23,6 +24,7 @@ describe("Telegram QA profiles", () => { const live = resolveTelegramQaScenarioIds({ providerMode: "live-frontier" }); const mock = resolveTelegramQaScenarioIds({ providerMode: "mock-openai" }); + expect(mock).toHaveLength(19); expect(live).not.toContain("telegram-long-final-reuses-preview"); expect(mock).toContain("telegram-long-final-reuses-preview"); expect(mock).not.toContain("telegram-assistant-transcript-role-boundary"); @@ -54,6 +56,13 @@ describe("Telegram QA profiles", () => { scenarioIds: ["telegram-startup-getme-live"], }), ).toThrow("execution.kind=flow"); + expect( + resolveTelegramQaScenarioIds({ + profile: "release", + providerMode: "mock-openai", + scenarioIds: ["channel-canary"], + }), + ).toEqual(["channel-canary"]); }); it("selects the native queue-validation regression as an explicit live scenario", () => { diff --git a/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.test.ts b/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.test.ts index 703b493f6050..8a5ac998b7e5 100644 --- a/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.test.ts +++ b/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.test.ts @@ -1,14 +1,30 @@ -import { describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const fetchWithSsrFGuard = vi.hoisted(() => vi.fn()); + +vi.mock("openclaw/plugin-sdk/ssrf-runtime", () => ({ fetchWithSsrFGuard })); + import { buildTelegramQaConfig, + callTelegramApi, isRecoverableTelegramQaPollError, normalizeTelegramObservedMessage, parseTelegramQaCredentialPayload, resolveTelegramQaRuntimeEnv, + TelegramQaApiError, + waitForTelegramPollRetryDelay, waitForTelegramChannelRunning, } from "./telegram-api.runtime.js"; describe("Telegram QA API boundary", () => { + beforeEach(() => { + fetchWithSsrFGuard.mockReset(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + it("parses env and leased credential payloads", () => { expect( resolveTelegramQaRuntimeEnv({ @@ -145,4 +161,137 @@ describe("Telegram QA API boundary", () => { expect(isRecoverableTelegramQaPollError(new Error("socket hang up"))).toBe(true); expect(isRecoverableTelegramQaPollError(new Error("Telegram unauthorized"))).toBe(false); }); + + it.each([ + { errorCode: 400, description: "Bad Request" }, + { errorCode: 401, description: "Unauthorized" }, + { errorCode: 404, description: "Not Found" }, + { + errorCode: 409, + description: + "Conflict: terminated by other getUpdates request; make sure that only one bot instance is running", + }, + ])("preserves typed terminal Telegram $errorCode errors", async ({ errorCode, description }) => { + const release = vi.fn(); + fetchWithSsrFGuard.mockResolvedValue({ + response: new Response( + JSON.stringify({ + ok: false, + error_code: errorCode, + description, + parameters: { retry_after: 3 }, + }), + { status: errorCode }, + ), + release, + }); + + const error = await callTelegramApi("placeholder", "getUpdates").catch( + (caught: unknown) => caught, + ); + + expect(error).toBeInstanceOf(TelegramQaApiError); + expect(error).toMatchObject({ + method: "getUpdates", + error_code: errorCode, + description, + parameters: { retry_after: 3 }, + status: errorCode, + }); + expect(isRecoverableTelegramQaPollError(error)).toBe(false); + expect(release).toHaveBeenCalledOnce(); + }); + + it.each([429, 500, 502])("retries typed transient Telegram %s errors", (errorCode) => { + expect( + isRecoverableTelegramQaPollError( + new TelegramQaApiError( + "getUpdates", + errorCode, + "transient Telegram failure", + undefined, + errorCode, + ), + ), + ).toBe(true); + }); + + it("honors retry_after and caps exponential poll backoff", async () => { + vi.useFakeTimers(); + const rateLimit = new TelegramQaApiError( + "getUpdates", + 429, + "Too Many Requests", + { retry_after: 3 }, + 429, + ); + const serverError = new TelegramQaApiError("getUpdates", 502, "Bad Gateway", undefined, 502); + const cases = [ + { error: rateLimit, attempt: 7, delayMs: 3_000 }, + ...[250, 500, 1_000, 2_000, 2_000].map((delayMs, index) => ({ + error: serverError, + attempt: index + 1, + delayMs, + })), + ...[250, 500, 1_000, 2_000, 2_000].map((delayMs, index) => ({ + error: new Error("fetch failed"), + attempt: index + 1, + delayMs, + })), + ]; + + for (const testCase of cases) { + let settled = false; + const waiting = waitForTelegramPollRetryDelay( + testCase.error, + testCase.attempt, + new AbortController().signal, + ).then(() => { + settled = true; + }); + await vi.advanceTimersByTimeAsync(testCase.delayMs - 1); + expect(settled).toBe(false); + await vi.advanceTimersByTimeAsync(1); + await waiting; + expect(settled).toBe(true); + } + }); + + it("aborts an in-flight Telegram poll retry delay", async () => { + const controller = new AbortController(); + const waiting = waitForTelegramPollRetryDelay( + new TelegramQaApiError("getUpdates", 429, "Too Many Requests", { retry_after: 60 }, 429), + 1, + controller.signal, + ); + + controller.abort(new Error("observer cleanup")); + + await expect(waiting).rejects.toThrow("aborted"); + }); + + it("preserves a non-JSON HTTP error as a typed Telegram status and releases transport", async () => { + const release = vi.fn(); + fetchWithSsrFGuard.mockResolvedValue({ + response: new Response("bad gateway", { + headers: { "content-type": "text/html" }, + status: 502, + }), + release, + }); + + const error = await callTelegramApi("placeholder", "getUpdates").catch( + (caught: unknown) => caught, + ); + + expect(error).toBeInstanceOf(TelegramQaApiError); + expect(error).toMatchObject({ + method: "getUpdates", + error_code: 502, + description: "getUpdates failed with status 502", + status: 502, + }); + expect(isRecoverableTelegramQaPollError(error)).toBe(true); + expect(release).toHaveBeenCalledOnce(); + }); }); diff --git a/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.ts b/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.ts index 642569465df5..f279c0bef1fe 100644 --- a/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.ts +++ b/extensions/qa-lab/src/live-transports/telegram/telegram-api.runtime.ts @@ -8,6 +8,11 @@ import { resolveTimerTimeoutMs, } from "openclaw/plugin-sdk/number-runtime"; import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; +import { + computeBackoff, + sleepWithAbort, + type BackoffPolicy, +} from "openclaw/plugin-sdk/runtime-env"; import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; import { isRecord, uniqueStrings } from "openclaw/plugin-sdk/string-coerce-runtime"; import { z } from "zod"; @@ -29,8 +34,26 @@ type TelegramApiEnvelope = { ok: boolean; result?: T; description?: string; + error_code?: number; + parameters?: unknown; }; +export class TelegramQaApiError extends Error { + override readonly name = "TelegramQaApiError"; + readonly ok = false; + + constructor( + readonly method: string, + readonly error_code: number, + readonly description: string, + readonly parameters: unknown, + readonly status: number, + options?: { cause?: unknown }, + ) { + super(`${method} failed (${error_code}): ${description}`, options); + } +} + type TelegramReplyMarkup = { inline_keyboard?: Array>; }; @@ -91,6 +114,12 @@ type TelegramGatewayClient = { }; const TELEGRAM_QA_DEFAULT_READY_TIMEOUT_MS = 45_000; +const TELEGRAM_QA_POLL_RETRY_BACKOFF: BackoffPolicy = { + initialMs: 250, + maxMs: 2_000, + factor: 2, + jitter: 0, +}; const TELEGRAM_QA_ENV_FIELDS = [ { field: "groupId", envKey: "OPENCLAW_QA_TELEGRAM_GROUP_ID" }, { field: "driverToken", envKey: "OPENCLAW_QA_TELEGRAM_DRIVER_BOT_TOKEN" }, @@ -336,14 +365,42 @@ export async function callTelegramApi( capture: false, }); try { - const payload = await readProviderJsonResponse>( - response, - `qa-lab-telegram-live.${method}`, - ); - if (!response.ok || !payload.ok || payload.result === undefined) { - throw new Error( - payload.description?.trim() || `${method} failed with status ${response.status}`, + let payload: TelegramApiEnvelope; + try { + const parsed = await readProviderJsonResponse( + response, + `qa-lab-telegram-live.${method}`, ); + if (!isRecord(parsed)) { + throw new Error(`qa-lab-telegram-live.${method}: malformed JSON response`); + } + payload = parsed as TelegramApiEnvelope; + } catch (error) { + if (!response.ok) { + throw new TelegramQaApiError( + method, + response.status, + `${method} failed with status ${response.status}`, + undefined, + response.status, + { cause: error }, + ); + } + throw error; + } + if (!response.ok || !payload.ok) { + const description = + payload.description?.trim() || `${method} failed with status ${response.status}`; + throw new TelegramQaApiError( + method, + typeof payload.error_code === "number" ? payload.error_code : response.status, + description, + payload.parameters, + response.status, + ); + } + if (payload.result === undefined) { + throw new Error(`${method} returned no result`); } return payload.result; } finally { @@ -352,6 +409,12 @@ export async function callTelegramApi( } export function isRecoverableTelegramQaPollError(error: unknown): boolean { + if (error && typeof error === "object" && "error_code" in error) { + const errorCode = (error as { error_code?: unknown }).error_code; + if (typeof errorCode === "number") { + return errorCode === 429 || errorCode >= 500; + } + } const message = formatErrorMessage(error).toLowerCase(); return [ "fetch failed", @@ -366,8 +429,31 @@ export function isRecoverableTelegramQaPollError(error: unknown): boolean { ].some((fragment) => message.includes(fragment)); } -export async function waitForTelegramPollRetryDelay(remainingMs = 250) { - await sleep(Math.min(250, Math.max(100, remainingMs))); +function readTelegramQaRetryAfterMs(error: unknown) { + if (!(error instanceof TelegramQaApiError) || error.error_code !== 429) { + return undefined; + } + const retryAfter = isRecord(error.parameters) ? error.parameters.retry_after : undefined; + return typeof retryAfter === "number" && Number.isFinite(retryAfter) && retryAfter >= 0 + ? retryAfter * 1_000 + : undefined; +} + +function resolveTelegramPollRetryDelayMs(error: unknown, attempt: number) { + return ( + readTelegramQaRetryAfterMs(error) ?? + computeBackoff(TELEGRAM_QA_POLL_RETRY_BACKOFF, Math.max(1, attempt)) + ); +} + +export async function waitForTelegramPollRetryDelay( + error: unknown, + attempt: number, + abortSignal: AbortSignal, +) { + await sleepWithAbort(resolveTelegramPollRetryDelayMs(error, attempt), abortSignal, { + ref: false, + }); } export async function flushTelegramUpdates(token: string) { diff --git a/extensions/qa-lab/src/providers/shared/mock-auth.ts b/extensions/qa-lab/src/providers/shared/mock-auth.ts index b7873623c76f..c27054966450 100644 --- a/extensions/qa-lab/src/providers/shared/mock-auth.ts +++ b/extensions/qa-lab/src/providers/shared/mock-auth.ts @@ -10,7 +10,7 @@ const QA_MOCK_AUTH_PROVIDERS = Object.freeze(["openai", "anthropic"] as const); /** Agent IDs the mock harness stages credentials under. */ const QA_MOCK_AUTH_AGENT_IDS = Object.freeze(["main", "qa"] as const); -function buildQaMockProfileId(provider: string): string { +export function buildQaMockProfileId(provider: string): string { return `qa-mock-${provider}`; } diff --git a/extensions/qa-lab/src/qa-transport.test.ts b/extensions/qa-lab/src/qa-transport.test.ts index e22287dc38ac..2eb95e07941f 100644 --- a/extensions/qa-lab/src/qa-transport.test.ts +++ b/extensions/qa-lab/src/qa-transport.test.ts @@ -94,6 +94,55 @@ describe("createQaStateBackedTransportAdapter", () => { expect(adapter.prepareFlow).toBeTypeOf("function"); expect(state.getSnapshot().messages).toHaveLength(0); }); + + it("adds redacted transport and bus-kind evidence to outbound timeouts", async () => { + const state = createQaBusState(); + state.addInboundMessage( + { + accountId: "sut", + conversation: { id: "private-chat-id", kind: "group" }, + senderId: "private-sender-id", + text: "private message content", + }, + "private-native-message-id", + ); + const adapter = createQaStateBackedTransportAdapter(state, { + id: "live", + label: "Live", + accountId: "sut", + requiredPluginIds: [], + supportedActions: [], + describeTransportState: () => + "telegram observer polls=3; updates=2; filtered=1; matched=1; update kinds=[message]; terminal error=none", + sendInbound: async (input) => state.addInboundMessage(input), + createGatewayConfig: () => ({}), + waitReady: async () => undefined, + buildAgentDelivery: ({ target }) => ({ + channel: "live", + to: target, + replyChannel: "live", + replyTo: target, + }), + handleAction: async () => undefined, + createReportNotes: () => [], + }); + + const error = await adapter + .waitForOutboundSequence?.({ + finalTextIncludes: "missing final", + timeoutMs: 5, + }) + .catch((caught: unknown) => caught); + const message = String(error); + + expect(message).toContain( + "telegram observer polls=3; updates=2; filtered=1; matched=1; update kinds=[message]; terminal error=none", + ); + expect(message).toContain("final bus-event kinds=[inbound-message]"); + expect(message).not.toMatch( + /private-chat-id|private-sender-id|private message content|private-native-message-id/u, + ); + }); }); describe("waitForQaTransportOutboundSequence", () => { @@ -136,6 +185,38 @@ describe("waitForQaTransportOutboundSequence", () => { }); }); + it("returns preview and final sends across distinct messages", async () => { + const state = createQaBusState(); + const preview = state.addOutboundMessage({ + accountId: "default", + text: "preview", + to: "dm:alice", + }); + const final = state.addOutboundMessage({ + accountId: "default", + text: "final marker", + to: "dm:alice", + }); + + const sequence = await waitForQaTransportOutboundSequence({ + accountId: "default", + input: { + conversationId: "alice", + finalSettleMs: 0, + finalTextIncludes: "final marker", + minimumPreviewEvents: 1, + timeoutMs: 100, + }, + readEvents: () => state.getSnapshot().events, + }); + + expect(sequence.events.map(({ kind, message }) => [kind, message.id])).toEqual([ + ["sent", preview.id], + ["sent", final.id], + ]); + expect(sequence.final).toMatchObject({ id: final.id, text: "final marker" }); + }); + it.each([ { description: "before the final marker", failureBeforeFinal: true }, { description: "after the final marker", failureBeforeFinal: false }, @@ -269,6 +350,13 @@ describe("waitForQaTransportOutboundSequence", () => { it("does not count an already-final send as a preview", async () => { const state = createQaBusState(); + state.addOutboundMessage({ + accountId: "default", + senderId: "openclaw", + text: "stale preview", + to: "dm:alice", + }); + const sinceCursor = state.getSnapshot().cursor; const final = state.addOutboundMessage({ accountId: "default", senderId: "openclaw", @@ -289,6 +377,7 @@ describe("waitForQaTransportOutboundSequence", () => { finalSettleMs: 0, finalTextIncludes: "final marker", minimumPreviewEvents: 1, + sinceCursor, timeoutMs: 20, }, readEvents: () => state.getSnapshot().events, diff --git a/extensions/qa-lab/src/qa-transport.ts b/extensions/qa-lab/src/qa-transport.ts index f4ed05bfe5c2..d2f4d748016f 100644 --- a/extensions/qa-lab/src/qa-transport.ts +++ b/extensions/qa-lab/src/qa-transport.ts @@ -193,6 +193,7 @@ export async function waitForQaTransportCondition( check: () => T | Promise | null | undefined, timeoutMs = 15_000, intervalMs = 100, + describeTimeout?: () => string, ): Promise { const pollIntervalMs = resolveTimerTimeoutMs(intervalMs, 100, 0); const startedAt = Date.now(); @@ -207,7 +208,8 @@ export async function waitForQaTransportCondition( } await sleep(Math.min(pollIntervalMs, remainingMs)); } - throw new Error(`timed out after ${timeoutMs}ms`); + const details = describeTimeout?.().trim(); + throw new Error(`timed out after ${timeoutMs}ms${details ? `; ${details}` : ""}`); } export function findFailureOutboundMessage( @@ -245,6 +247,7 @@ function createFailureAwareTransportWaitForCondition(state: QaTransportState, ac check: () => T | Promise | null | undefined, timeoutMs = 15_000, intervalMs = 100, + describeTimeout?: () => string, ): Promise { const sinceIndex = state.getSnapshot().messages.length; return await waitForQaTransportCondition( @@ -264,10 +267,37 @@ function createFailureAwareTransportWaitForCondition(state: QaTransportState, ac }, timeoutMs, intervalMs, + describeTimeout, ); }; } +const QA_TRANSPORT_TIMEOUT_EVENT_LIMIT = 8; + +function describeQaTransportTimeout(params: { + accountId: string; + describeTransportState?: () => string; + state: QaTransportState; +}) { + const eventKinds = params.state + .getSnapshot() + .events.filter((event) => event.accountId === params.accountId) + .slice(-QA_TRANSPORT_TIMEOUT_EVENT_LIMIT) + .map((event) => event.kind); + let transportState: string | undefined; + try { + transportState = params.describeTransportState?.().trim(); + } catch { + transportState = "transport state unavailable"; + } + return [ + transportState, + `final bus-event kinds=[${eventKinds.length > 0 ? eventKinds.join(",") : "none"}]`, + ] + .filter(Boolean) + .join("; "); +} + type QaTransportAdapterDefinition = Awaited< ReturnType["create"]> >; @@ -296,6 +326,7 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte readonly state: QaTransportState; readonly waitForCondition: QaTransportAdapter["waitForCondition"]; private readonly assertTransportHealthy: () => void; + private readonly describeTimeout: () => string; constructor(params: { id: string; @@ -305,6 +336,8 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte supportedActions?: readonly QaTransportActionName[]; state: QaTransportState; assertTransportHealthy?: () => void; + describeTimeout?: () => string; + describeTransportState?: () => string; }) { this.id = params.id; this.label = params.label; @@ -313,6 +346,14 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte this.supportedActions = params.supportedActions ?? []; this.state = params.state; this.assertTransportHealthy = params.assertTransportHealthy ?? (() => undefined); + this.describeTimeout = + params.describeTimeout ?? + (() => + describeQaTransportTimeout({ + accountId: this.accountId, + describeTransportState: params.describeTransportState, + state: this.state, + })); const waitForCondition = createFailureAwareTransportWaitForCondition( this.state, this.accountId, @@ -325,6 +366,7 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte }, timeoutMs, intervalMs, + this.describeTimeout, ); } @@ -375,32 +417,37 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte } async waitForOutbound(input: QaTransportOutboundMatch) { - return await waitForQaTransportCondition(() => { - this.assertTransportHealthy(); - assertNoFailureReplies(this.state, { - accountId: this.accountId, - sinceIndex: input.sinceIndex, - cursorSpace: "outbound", - }); - return this.outboundSince(input.sinceIndex).find((message) => { - if (message.deleted) { - return false; - } - if (input.conversation && message.conversation.id !== input.conversation.id) { - return false; - } - if (input.conversation && message.conversation.kind !== input.conversation.kind) { - return false; - } - if (input.senderId && message.senderId !== input.senderId) { - return false; - } - if (input.threadId && message.threadId !== input.threadId) { - return false; - } - return !input.textIncludes || message.text.includes(input.textIncludes); - }); - }, input.timeoutMs); + return await waitForQaTransportCondition( + () => { + this.assertTransportHealthy(); + assertNoFailureReplies(this.state, { + accountId: this.accountId, + sinceIndex: input.sinceIndex, + cursorSpace: "outbound", + }); + return this.outboundSince(input.sinceIndex).find((message) => { + if (message.deleted) { + return false; + } + if (input.conversation && message.conversation.id !== input.conversation.id) { + return false; + } + if (input.conversation && message.conversation.kind !== input.conversation.kind) { + return false; + } + if (input.senderId && message.senderId !== input.senderId) { + return false; + } + if (input.threadId && message.threadId !== input.threadId) { + return false; + } + return !input.textIncludes || message.text.includes(input.textIncludes); + }); + }, + input.timeoutMs, + undefined, + this.describeTimeout, + ); } private outboundSince(sinceIndex = 0) { @@ -416,6 +463,12 @@ export function createQaStateBackedTransportAdapter( state: QaTransportState, params: QaTransportAdapterDefinition, ): QaTransportAdapter { + const describeTimeout = () => + describeQaTransportTimeout({ + accountId: params.accountId, + describeTransportState: params.describeTransportState, + state, + }); const adapter = new (class extends QaStateBackedTransportAdapter { createGatewayConfig = params.createGatewayConfig; waitReady = params.waitReady; @@ -437,6 +490,7 @@ export function createQaStateBackedTransportAdapter( supportedActions: params.supportedActions, state, assertTransportHealthy: params.assertTransportHealthy, + describeTimeout, }); Object.assign(adapter, { ...(params.sendNativeCommand ? { sendNativeCommand: params.sendNativeCommand } : {}), @@ -450,6 +504,7 @@ export function createQaStateBackedTransportAdapter( params.assertTransportHealthy?.(); return state.getSnapshot().events; }, + describeTimeout, })), ...(params.createRuntimeEnvPatch ? { createRuntimeEnvPatch: params.createRuntimeEnvPatch } @@ -488,79 +543,95 @@ export async function waitForQaTransportOutboundSequence(params: { readEvents: () => | readonly (QaBusEvent | QaTransportOutboundEvent)[] | Promise; + describeTimeout?: () => string; }): Promise { const minimumPreviewEvents = params.input.minimumPreviewEvents ?? 1; const finalSettleMs = params.input.finalSettleMs ?? 300; let stableCursor: number | null = null; let stableSince = 0; - return await waitForQaTransportCondition(async () => { - const ownedEvents = (await params.readEvents()) - .filter((event) => event.cursor > (params.input.sinceCursor ?? 0)) - .map((event) => - isQaTransportOutboundEvent(event) ? event : normalizeQaBusOutboundEvent(event), - ) - .filter((event): event is QaTransportOutboundEvent => event !== null) - .filter( - ({ message }) => message.accountId === params.accountId && message.direction === "outbound", + return await waitForQaTransportCondition( + async () => { + const ownedEvents = (await params.readEvents()) + .filter((event) => event.cursor > (params.input.sinceCursor ?? 0)) + .map((event) => + isQaTransportOutboundEvent(event) ? event : normalizeQaBusOutboundEvent(event), + ) + .filter((event): event is QaTransportOutboundEvent => event !== null) + .filter( + ({ message }) => + message.accountId === params.accountId && message.direction === "outbound", + ); + // Failures belong to the account, even when a different conversation has a matching final. + for (const { kind, message } of ownedEvents) { + const failureReply = + kind === "deleted" || message.deleted + ? undefined + : extractQaFailureReplyText(message.text); + if (failureReply) { + throw new Error(failureReply); + } + } + const events = ownedEvents.filter(({ message }) => { + if ( + params.input.conversationId && + message.conversation.id !== params.input.conversationId + ) { + return false; + } + return !params.input.threadId || message.threadId === params.input.threadId; + }); + const finalIndex = events.findLastIndex( + ({ kind, message }) => + kind !== "deleted" && + !message.deleted && + message.text.includes(params.input.finalTextIncludes), ); - // Failures belong to the account, even when a different conversation has a matching final. - for (const { kind, message } of ownedEvents) { - const failureReply = - kind === "deleted" || message.deleted ? undefined : extractQaFailureReplyText(message.text); - if (failureReply) { - throw new Error(failureReply); + if (finalIndex < 0) { + return undefined; } - } - const events = ownedEvents.filter(({ message }) => { - if (params.input.conversationId && message.conversation.id !== params.input.conversationId) { - return false; + const candidate = events[finalIndex]; + if (!candidate) { + return undefined; } - return !params.input.threadId || message.threadId === params.input.threadId; - }); - const finalIndex = events.findLastIndex( - ({ kind, message }) => - kind !== "deleted" && - !message.deleted && - message.text.includes(params.input.finalTextIncludes), - ); - if (finalIndex < 0) { - return undefined; - } - const candidate = events[finalIndex]; - if (!candidate) { - return undefined; - } - const sequenceEvents = events.filter(({ message }) => message.id === candidate.message.id); - const latest = sequenceEvents.at(-1); - if ( - !latest || - latest.kind === "deleted" || - latest.message.deleted || - !latest.message.text.includes(params.input.finalTextIncludes) - ) { - stableCursor = null; - return undefined; - } - const previewEvents = sequenceEvents.filter( - ({ cursor, kind, message }) => - cursor < candidate.cursor && - kind !== "deleted" && - !message.text.includes(params.input.finalTextIncludes), - ); - if (previewEvents.length < minimumPreviewEvents) { - return undefined; - } - if (stableCursor !== latest.cursor) { - stableCursor = latest.cursor; - stableSince = Date.now(); - return finalSettleMs === 0 ? { events: sequenceEvents, final: latest.message } : undefined; - } - if (Date.now() - stableSince < finalSettleMs) { - return undefined; - } - return { - events: sequenceEvents, - final: latest.message, - }; - }, params.input.timeoutMs); + const finalLineage = events.filter(({ message }) => message.id === candidate.message.id); + const latest = finalLineage.at(-1); + if ( + !latest || + latest.kind === "deleted" || + latest.message.deleted || + !latest.message.text.includes(params.input.finalTextIncludes) + ) { + stableCursor = null; + return undefined; + } + const previewEvents = events.filter( + ({ cursor, kind, message }) => + cursor < candidate.cursor && + kind !== "deleted" && + !message.text.includes(params.input.finalTextIncludes), + ); + if (previewEvents.length < minimumPreviewEvents) { + return undefined; + } + const sequenceCursors = new Set( + [...previewEvents, ...finalLineage].map(({ cursor }) => cursor), + ); + const sequenceEvents = events.filter(({ cursor }) => sequenceCursors.has(cursor)); + if (stableCursor !== latest.cursor) { + stableCursor = latest.cursor; + stableSince = Date.now(); + return finalSettleMs === 0 ? { events: sequenceEvents, final: latest.message } : undefined; + } + if (Date.now() - stableSince < finalSettleMs) { + return undefined; + } + return { + events: sequenceEvents, + final: latest.message, + }; + }, + params.input.timeoutMs, + undefined, + params.describeTimeout, + ); } diff --git a/extensions/qa-lab/src/scenario-catalog-channels.test.ts b/extensions/qa-lab/src/scenario-catalog-channels.test.ts index 8c8a67219ebd..b81fc546418e 100644 --- a/extensions/qa-lab/src/scenario-catalog-channels.test.ts +++ b/extensions/qa-lab/src/scenario-catalog-channels.test.ts @@ -204,6 +204,7 @@ describe("qa scenario catalog channel contracts", () => { expect(scenario.execution.channel).toBeUndefined(); expect(scenario.execution.channels).toEqual(["qa-channel", "telegram"]); + expect(scenario.execution.retryCount).toBe(0); expect(scenario.coverage?.primary).toEqual(["channels.streaming-final-reply"]); expect(scenario.coverage?.secondary).toEqual([`${agentRuntime}.streaming-replies-delivery`]); expect(scenario.gatewayConfigPatch).toMatchObject({ @@ -212,6 +213,12 @@ describe("qa scenario catalog channel contracts", () => { expect(scenario.gatewayConfigPatch).not.toHaveProperty("channels.telegram.groups"); }); + it("keeps the shared channel canary eligible for QA Channel and Telegram", () => { + const scenario = requireFlowScenario(readQaScenarioById("channel-canary")); + + expect(scenario.execution.channels).toEqual(["qa-channel", "telegram"]); + }); + it("keeps transcript-role delivery on the Crabline driver", () => { const scenario = readQaScenarioById("telegram-assistant-transcript-role-boundary"); const config = readQaScenarioExecutionConfig("telegram-assistant-transcript-role-boundary") as diff --git a/qa/scenarios/channels/channel-canary.yaml b/qa/scenarios/channels/channel-canary.yaml index 1f98ebb3c5b1..56bb62870d65 100644 --- a/qa/scenarios/channels/channel-canary.yaml +++ b/qa/scenarios/channels/channel-canary.yaml @@ -21,6 +21,9 @@ scenario: - extensions/qa-lab/src/suite.ts execution: kind: flow + channels: + - qa-channel + - telegram summary: Run the shared channel canary through QA Channel, Crabline, or a live adapter. transportPolicy: requireGroupMention: true diff --git a/qa/scenarios/channels/channel-message-flows.yaml b/qa/scenarios/channels/channel-message-flows.yaml index c7fd1dfa0e07..ae5786fd12b3 100644 --- a/qa/scenarios/channels/channel-message-flows.yaml +++ b/qa/scenarios/channels/channel-message-flows.yaml @@ -29,6 +29,7 @@ scenario: - extensions/telegram/src/draft-stream.ts execution: kind: flow + retryCount: 0 channels: - qa-channel - telegram diff --git a/scripts/e2e/lib/npm-telegram-live/prepare-package.mjs b/scripts/e2e/lib/npm-telegram-live/prepare-package.mjs index e284d6f27d84..2abb380ff562 100644 --- a/scripts/e2e/lib/npm-telegram-live/prepare-package.mjs +++ b/scripts/e2e/lib/npm-telegram-live/prepare-package.mjs @@ -1,26 +1,24 @@ -// Prepares package manifests for npm Telegram live E2E scenarios. +// Prepares the trusted harness manifest for npm Telegram live E2E scenarios. import fs from "node:fs"; +import { privateLocalOnlyPluginSdkEntrypoints } from "../../../lib/plugin-sdk-entries.mjs"; const packageJsonPaths = process.argv.slice(2); -if (packageJsonPaths.length !== 2) { - throw new Error( - `expected exactly two ephemeral package manifests, got ${packageJsonPaths.length}`, - ); +if (packageJsonPaths.length !== 1) { + throw new Error("expected exactly one trusted harness package.json path"); } -for (const packageJsonPath of packageJsonPaths) { - const pkg = JSON.parse(fs.readFileSync(packageJsonPath, "utf8")); - pkg.exports = pkg.exports && typeof pkg.exports === "object" ? pkg.exports : {}; - if (!pkg.exports["./plugin-sdk/gateway-runtime"]) { - pkg.exports["./plugin-sdk/gateway-runtime"] = { - types: "./dist/plugin-sdk/gateway-runtime.d.ts", - default: "./dist/plugin-sdk/gateway-runtime.js", +const packageJsonPath = packageJsonPaths[0]; +const pkg = JSON.parse(fs.readFileSync(packageJsonPath, "utf8")); +pkg.exports = pkg.exports && typeof pkg.exports === "object" ? pkg.exports : {}; + +// Private QA builds emit these two harness-only facades outside the regular SDK inventory. +for (const subpath of [...privateLocalOnlyPluginSdkEntrypoints, "qa-lab", "qa-runtime"]) { + const exportPath = `./plugin-sdk/${subpath}`; + if (!pkg.exports[exportPath]) { + pkg.exports[exportPath] = { + default: `./dist/plugin-sdk/${subpath}.js`, }; } - if (!pkg.exports["./plugin-sdk/qa-runtime"]) { - pkg.exports["./plugin-sdk/qa-runtime"] = { - default: "./.openclaw-qa-harness-dist/plugin-sdk/qa-runtime.js", - }; - } - fs.writeFileSync(packageJsonPath, `${JSON.stringify(pkg, null, 2)}\n`); } + +fs.writeFileSync(packageJsonPath, `${JSON.stringify(pkg, null, 2)}\n`); diff --git a/scripts/e2e/npm-telegram-live-docker.sh b/scripts/e2e/npm-telegram-live-docker.sh index 9fccc88d1dd2..d50f8d51ed7b 100755 --- a/scripts/e2e/npm-telegram-live-docker.sh +++ b/scripts/e2e/npm-telegram-live-docker.sh @@ -213,11 +213,16 @@ docker_e2e_build_or_reuse "$IMAGE_NAME" npm-telegram-live "$ROOT_DIR/scripts/e2e mkdir -p "$ROOT_DIR/.artifacts/qa-e2e" mkdir -p "$OUTPUT_DIR_HOST" npm_prefix_host="$(mktemp -d "$ROOT_DIR/.artifacts/qa-e2e/npm-telegram-live-prefix.XXXXXX")" +harness_root="$(mktemp -d "$ROOT_DIR/.artifacts/qa-e2e/npm-telegram-live-harness.XXXXXX")" +harness_package_json="$harness_root/package.json" +cp "$ROOT_DIR/package.json" "$harness_package_json" +node "$ROOT_DIR/scripts/e2e/lib/npm-telegram-live/prepare-package.mjs" "$harness_package_json" cleanup() { local rc=$? trap - EXIT printf 'schema=1\nexit_code=%s\nlive_output=job_log\n' "$rc" > "$OUTPUT_DIR_HOST/run-metadata.txt" rm -rf "$npm_prefix_host" + rm -rf "$harness_root" exit "$rc" } trap cleanup EXIT @@ -408,14 +413,18 @@ command -v openclaw openclaw --version EOF -# Mount only QA harness source; the SUT itself, including bundled plugin runtime, -# is the installed package candidate. +# Mount the trusted current-source QA harness separately from the installed +# package candidate. The candidate remains the absolute CLI/runtime SUT. run_logged_print_heartbeat "npm-telegram-live-suite" 60 docker_e2e_run_with_harness \ "${docker_env[@]}" \ -v "$ROOT_DIR/.artifacts:/app/.artifacts" \ -v "$OUTPUT_DIR_HOST:$OUTPUT_DIR_CONTAINER" \ - -v "$ROOT_DIR/dist:/app/.openclaw-qa-harness-dist:ro" \ - -v "$ROOT_DIR/extensions/qa-lab:/app/extensions/qa-lab:ro" \ + -v "$harness_package_json:/app/package.json:ro" \ + -v "$ROOT_DIR/dist:/app/dist:ro" \ + -v "$ROOT_DIR/node_modules:/trusted-harness/node_modules:ro" \ + -v "$ROOT_DIR/packages:/app/packages:ro" \ + -v "$ROOT_DIR/extensions:/app/extensions:ro" \ + -v "$ROOT_DIR/taxonomy.yaml:/app/taxonomy.yaml:ro" \ -v "$ROOT_DIR/qa/scenarios:/app/qa/scenarios:ro" \ -v "$ROOT_DIR/taxonomy.yaml:/app/taxonomy.yaml:ro" \ -v "$npm_prefix_host:/npm-global" \ @@ -428,6 +437,7 @@ export HOME="$runtime_home" export NPM_CONFIG_PREFIX="/npm-global" export PATH="$NPM_CONFIG_PREFIX/bin:$PATH" export OPENCLAW_NPM_TELEGRAM_REPO_ROOT="/app" +sut_command="/npm-global/bin/openclaw" dump_hotpath_logs() { local status="$1" @@ -445,70 +455,53 @@ dump_hotpath_logs() { } trap 'status=$?; dump_hotpath_logs "$status"; exit "$status"' ERR -command -v openclaw -openclaw_e2e_run_command openclaw --version +test -x "$sut_command" +openclaw_e2e_run_command "$sut_command" --version mkdir -p /app/node_modules -openclaw_package_dir="/npm-global/lib/node_modules/openclaw" -# The mounted QA harness imports openclaw/plugin-sdk and package dependencies; -# point those imports at the installed package without copying source plugins into the test image. -rm -rf /app/node_modules/openclaw -ln -sfnT "$openclaw_package_dir" /app/node_modules/openclaw -rm -rf /app/dist -ln -sfnT "$openclaw_package_dir/dist" /app/dist -rm -rf "$openclaw_package_dir/.openclaw-qa-harness-dist" -ln -sfnT /app/.openclaw-qa-harness-dist "$openclaw_package_dir/.openclaw-qa-harness-dist" -cp "$openclaw_package_dir/package.json" /app/package.json -node scripts/e2e/lib/npm-telegram-live/prepare-package.mjs \ - /app/package.json \ - /app/node_modules/openclaw/package.json -for deps_dir in "$openclaw_package_dir/node_modules" /npm-global/lib/node_modules; do - [ -d "$deps_dir" ] || continue - for dependency_dir in "$deps_dir"/*; do - [ -e "$dependency_dir" ] || continue - dependency_name="$(basename "$dependency_dir")" - case "$dependency_name" in - .bin | openclaw) - continue - ;; - @*) - [ -d "$dependency_dir" ] || continue - mkdir -p "/app/node_modules/$dependency_name" - for scoped_dependency_dir in "$dependency_dir"/*; do - [ -e "$scoped_dependency_dir" ] || continue - scoped_dependency_name="$(basename "$scoped_dependency_dir")" - rm -rf "/app/node_modules/$dependency_name/$scoped_dependency_name" - ln -sfnT "$scoped_dependency_dir" "/app/node_modules/$dependency_name/$scoped_dependency_name" - done - ;; - *) - rm -rf "/app/node_modules/$dependency_name" - ln -sfnT "$dependency_dir" "/app/node_modules/$dependency_name" - ;; - esac - done -done - -link_installed_package_dependency() { - local name="$1" - local source="/npm-global/lib/node_modules/openclaw/node_modules/$name" +link_harness_dependency() { + local source="$1" + local name="$2" local target="/app/node_modules/$name" - if [ ! -e "$source" ]; then - echo "Installed package dependency is missing: $name" >&2 - return 1 - fi mkdir -p "$(dirname "$target")" - ln -sfn "$source" "$target" + ln -sfnT "$source" "$target" } -# QA Lab is intentionally mounted as harness source, so its package-local -# runtime imports must resolve from the installed package dependency tree. -for dependency in \ - @modelcontextprotocol/sdk \ - yaml \ - zod; do - link_installed_package_dependency "$dependency" +# External dependencies resolve from the trusted install, not the candidate. +for dependency_dir in /trusted-harness/node_modules/* /trusted-harness/node_modules/.[!.]*; do + [ -e "$dependency_dir" ] || continue + dependency_name="$(basename "$dependency_dir")" + case "$dependency_name" in + .bin | openclaw) + continue + ;; + @*) + [ -d "$dependency_dir" ] || continue + for scoped_dependency_dir in "$dependency_dir"/*; do + [ -e "$scoped_dependency_dir" ] || continue + scoped_dependency_name="$(basename "$scoped_dependency_dir")" + link_harness_dependency \ + "$scoped_dependency_dir" \ + "$dependency_name/$scoped_dependency_name" + done + ;; + *) + link_harness_dependency "$dependency_dir" "$dependency_name" + ;; + esac done +# Workspace links must resolve under /app even when pnpm linked them relative to +# the checkout path used by the workflow. +for workspace_dir in /app/packages/* /app/extensions/*; do + [ -f "$workspace_dir/package.json" ] || continue + workspace_name="$(node -e \ + 'const fs = require("node:fs"); const pkg = JSON.parse(fs.readFileSync(process.argv[1], "utf8")); process.stdout.write(pkg.name || "");' \ + "$workspace_dir/package.json")" + [ -n "$workspace_name" ] || continue + link_harness_dependency "$workspace_dir" "$workspace_name" +done +link_harness_dependency /app openclaw + if [ "${OPENCLAW_NPM_TELEGRAM_SKIP_HOTPATH:-0}" != "1" ]; then hotpath_home="$(mktemp -d "/tmp/openclaw-npm-telegram-hotpath.XXXXXX")" export HOME="$hotpath_home" @@ -519,7 +512,7 @@ if [ "${OPENCLAW_NPM_TELEGRAM_SKIP_HOTPATH:-0}" != "1" ]; then hotpath_model_value="$OPENAI_API_KEY" fi hotpath_channel_value="$(printf '%s:%s' 123456 "$hotpath_placeholder")" - OPENAI_API_KEY="$hotpath_model_value" openclaw_e2e_run_command openclaw onboard \ + OPENAI_API_KEY="$hotpath_model_value" openclaw_e2e_run_command "$sut_command" onboard \ --non-interactive --accept-risk \ --mode local \ --auth-choice openai-api-key \ @@ -532,13 +525,13 @@ if [ "${OPENCLAW_NPM_TELEGRAM_SKIP_HOTPATH:-0}" != "1" ]; then --skip-health \ --json >/tmp/openclaw-npm-telegram-onboard.json /tmp/openclaw-npm-telegram-channel-add.log 2>&1 /tmp/openclaw-npm-telegram-doctor-fix.log 2>&1 /tmp/openclaw-npm-telegram-doctor-check.log 2>&1 /tmp/openclaw-npm-telegram-channel-add.log 2>&1 /tmp/openclaw-npm-telegram-doctor-fix.log 2>&1 /tmp/openclaw-npm-telegram-doctor-check.log 2>&1 , +) { + if (!options) { + return [...scenarioIds]; + } + return [ + options.scenarioId, + ...scenarioIds.filter((scenarioId) => scenarioId !== options.scenarioId), + ]; +} + async function shouldFailPackageTelegramRun( result: { summaryPath: string }, env: NodeJS.ProcessEnv = process.env, @@ -141,8 +154,15 @@ async function resolveTrustedOpenClawCommand( } async function main() { - const { runQaTelegramSuite } = - await import("../../extensions/qa-lab/src/live-transports/telegram/cli.runtime.ts"); + const [ + { runQaTelegramSuite }, + { resolveTelegramQaScenarioIds }, + { DEFAULT_QA_LIVE_PROVIDER_MODE }, + ] = await Promise.all([ + import("../../extensions/qa-lab/src/live-transports/telegram/cli.runtime.ts"), + import("../../extensions/qa-lab/src/live-transports/telegram/scenario-selection.ts"), + import("../../extensions/qa-lab/src/providers/index.ts"), + ]); const rawSutOpenClawCommand = process.env.OPENCLAW_NPM_TELEGRAM_SUT_COMMAND?.trim(); if (!rawSutOpenClawCommand) { throw new Error("Missing OPENCLAW_NPM_TELEGRAM_SUT_COMMAND."); @@ -152,18 +172,29 @@ async function main() { const repoRoot = path.resolve(process.env.OPENCLAW_NPM_TELEGRAM_REPO_ROOT ?? process.cwd()); const outputDir = resolvePackageTelegramOutputDir(process.env, repoRoot); const scenarioIds = splitCsv(process.env.OPENCLAW_NPM_TELEGRAM_SCENARIOS); + const providerMode = + (process.env.OPENCLAW_NPM_TELEGRAM_PROVIDER_MODE as QaProviderMode | undefined) ?? + DEFAULT_QA_LIVE_PROVIDER_MODE; + const primaryModel = process.env.OPENCLAW_NPM_TELEGRAM_MODEL; + const resolvedScenarioIds = resolveTelegramQaScenarioIds({ + providerMode, + primaryModel, + scenarioIds, + }); + const rttOptions = resolveRttOptions(process.env, scenarioIds); const result = await runQaTelegramSuite({ allowFailures: true, failFast: true, repoRoot, outputDir, sutOpenClawCommand, - providerMode: process.env.OPENCLAW_NPM_TELEGRAM_PROVIDER_MODE as QaProviderMode | undefined, - primaryModel: process.env.OPENCLAW_NPM_TELEGRAM_MODEL, + providerMode, + primaryModel, alternateModel: process.env.OPENCLAW_NPM_TELEGRAM_ALT_MODEL, fastMode: parseBoolean(process.env.OPENCLAW_NPM_TELEGRAM_FAST), scenarioIds, - roundTripProbe: createRoundTripProbe(resolveRttOptions(process.env, scenarioIds)), + resolvedScenarioIds: prioritizeRoundTripProbeScenario(resolvedScenarioIds, rttOptions), + roundTripProbe: createRoundTripProbe(rttOptions), sutAccountId: process.env.OPENCLAW_NPM_TELEGRAM_SUT_ACCOUNT, credentialSource: resolveCredentialSource(process.env), credentialRole: resolveCredentialRole(process.env), @@ -209,6 +240,7 @@ export const testing = { resolveCredentialRole, resolveCredentialSource, createRoundTripProbe, + prioritizeRoundTripProbeScenario, resolveRttOptions, resolveTrustedOpenClawCommand, shouldFailPackageTelegramRun, diff --git a/src/plugin-sdk/qa-runner-runtime.ts b/src/plugin-sdk/qa-runner-runtime.ts index 622ef0c53798..2750510efd5d 100644 --- a/src/plugin-sdk/qa-runner-runtime.ts +++ b/src/plugin-sdk/qa-runner-runtime.ts @@ -121,6 +121,7 @@ type QaRunnerTransportAdapterDefinition = { requiredPluginIds: readonly string[]; supportedActions: readonly ("delete" | "edit" | "react" | "thread-create")[]; assertTransportHealthy?: () => void; + describeTransportState?: () => string; resetTransport?: () => void | Promise; sendInbound: (input: QaBusInboundMessageInput) => Promise; sendNativeCommand?: ( diff --git a/test/scripts/lint-suppressions.test.ts b/test/scripts/lint-suppressions.test.ts index 15f5c486071e..8181549ebc9e 100644 --- a/test/scripts/lint-suppressions.test.ts +++ b/test/scripts/lint-suppressions.test.ts @@ -196,6 +196,7 @@ describe("production lint suppressions", () => { "extensions/discord/src/test-support/provider.test-support.ts|typescript/no-unnecessary-type-parameters|1", "extensions/feishu/src/bitable.ts|typescript/no-unnecessary-type-parameters|1", "extensions/matrix/src/onboarding.test-harness.ts|typescript/no-unnecessary-type-parameters|1", + "extensions/qa-lab/src/gateway-child.ts|preserve-caught-error|1", "extensions/slack/src/monitor/provider-support.ts|typescript/no-unnecessary-type-parameters|1", "scripts/changed-lanes.mjs|typescript/no-base-to-string|2", "scripts/changed-lanes.mjs|typescript/restrict-template-expressions|2", diff --git a/test/scripts/npm-telegram-live.test.ts b/test/scripts/npm-telegram-live.test.ts index 9700d80982a1..5f936aaa78ce 100644 --- a/test/scripts/npm-telegram-live.test.ts +++ b/test/scripts/npm-telegram-live.test.ts @@ -1,11 +1,12 @@ // Npm Telegram Live tests cover npm telegram live script behavior. -import { spawnSync } from "node:child_process"; +import { execFileSync } from "node:child_process"; import { mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import path from "node:path"; import { fileURLToPath } from "node:url"; import { afterEach, describe, expect, it } from "vitest"; import { testing } from "../../scripts/e2e/npm-telegram-live-runner.ts"; +import { privateLocalOnlyPluginSdkEntrypoints } from "../../scripts/lib/plugin-sdk-entries.mjs"; const TEST_DIR = path.dirname(fileURLToPath(import.meta.url)); const DOCKER_SCRIPT_PATH = path.resolve(TEST_DIR, "../../scripts/e2e/npm-telegram-live-docker.sh"); @@ -51,7 +52,7 @@ describe("package Telegram live Docker E2E", () => { it("installs the package candidate before forwarding runtime secrets", () => { const script = readFileSync(DOCKER_SCRIPT_PATH, "utf8"); const installRunStart = script.indexOf('echo "Running package Telegram live Docker E2E'); - const installRunEnd = script.indexOf("# Mount only QA harness source"); + const installRunEnd = script.indexOf("# Mount the trusted current-source QA harness"); const installRun = script.slice(installRunStart, installRunEnd); expect(installRunStart).toBeGreaterThanOrEqual(0); @@ -95,7 +96,7 @@ describe("package Telegram live Docker E2E", () => { it("bounds installed-package hot path OpenClaw commands", () => { const script = readFileSync(DOCKER_SCRIPT_PATH, "utf8"); - const runtimeRunStart = script.indexOf("# Mount only QA harness source"); + const runtimeRunStart = script.indexOf("# Mount the trusted current-source QA harness"); const runtimeRun = script.slice(runtimeRunStart); expect(runtimeRunStart).toBeGreaterThanOrEqual(0); @@ -103,15 +104,19 @@ describe("package Telegram live Docker E2E", () => { '-e OPENCLAW_E2E_COMMAND_TIMEOUT="${OPENCLAW_E2E_COMMAND_TIMEOUT:-300s}"', ); expect(runtimeRun).toContain("source scripts/lib/openclaw-e2e-instance.sh"); - expect(runtimeRun).toContain("openclaw_e2e_run_command openclaw --version"); - expect(runtimeRun).toContain("openclaw_e2e_run_command openclaw onboard"); + expect(runtimeRun).toContain('sut_command="/npm-global/bin/openclaw"'); + expect(runtimeRun).toContain('openclaw_e2e_run_command "$sut_command" --version'); + expect(runtimeRun).toContain('openclaw_e2e_run_command "$sut_command" onboard'); expect(runtimeRun).toContain( - 'OPENAI_API_KEY="$hotpath_model_value" openclaw_e2e_run_command openclaw onboard', + 'OPENAI_API_KEY="$hotpath_model_value" openclaw_e2e_run_command "$sut_command" onboard', ); expect(runtimeRun).not.toContain("export OPENAI_API_KEY="); - expect(runtimeRun).toContain("openclaw_e2e_run_command openclaw channels add"); - expect(runtimeRun).toContain("openclaw_e2e_run_command openclaw doctor --fix"); - expect(runtimeRun).toContain("openclaw_e2e_run_command openclaw doctor --non-interactive"); + expect(runtimeRun).toContain('openclaw_e2e_run_command "$sut_command" channels add'); + expect(runtimeRun).toContain('openclaw_e2e_run_command "$sut_command" doctor --fix'); + expect(runtimeRun).toContain( + 'openclaw_e2e_run_command "$sut_command" doctor --non-interactive', + ); + expect(runtimeRun).toContain('export OPENCLAW_NPM_TELEGRAM_SUT_COMMAND="$sut_command"'); expect(runtimeRun).toContain('openclaw_e2e_print_log "$file"'); expect(runtimeRun).not.toContain("sed -n '1,220p'"); expect(runtimeRun).not.toMatch(/^\s*openclaw (onboard|channels add|doctor )/mu); @@ -215,8 +220,11 @@ describe("package Telegram live Docker E2E", () => { it("keeps the installed OpenClaw command as the package SUT", async () => { const prefix = mkTempRoot(); const command = path.join(prefix, "bin", "openclaw"); + const harnessCommand = path.join(mkTempRoot(), "bin", "openclaw"); mkdirSync(path.dirname(command), { recursive: true }); - writeFileSync(command, "#!/bin/sh\n"); + mkdirSync(path.dirname(harnessCommand), { recursive: true }); + writeFileSync(command, "#!/bin/sh\n", { mode: 0o755 }); + writeFileSync(harnessCommand, "#!/bin/sh\n", { mode: 0o755 }); await expect( testing.resolveTrustedOpenClawCommand(command, { @@ -226,6 +234,11 @@ describe("package Telegram live Docker E2E", () => { executablePath: command, usePackagedPlugins: true, }); + await expect( + testing.resolveTrustedOpenClawCommand(harnessCommand, { + NPM_CONFIG_PREFIX: prefix, + }), + ).rejects.toThrow("OPENCLAW_NPM_TELEGRAM_SUT_COMMAND must resolve inside NPM_CONFIG_PREFIX."); }); it("mounts the QA taxonomy without exposing the repository root", () => { @@ -275,94 +288,69 @@ describe("package Telegram live Docker E2E", () => { expect(script).toContain("OPENCLAW_NPM_TELEGRAM_RTT_CHECKS"); }); - it("keeps candidate runtime authoritative while mounting private QA dist separately", () => { + it("isolates the trusted private QA harness from the installed package candidate", () => { const script = readFileSync(DOCKER_SCRIPT_PATH, "utf8"); - const preparePackage = readFileSync(PREPARE_PACKAGE_PATH, "utf8"); - const gatewayRpcClient = readFileSync( - path.resolve(TEST_DIR, "../../extensions/qa-lab/src/gateway-rpc-client.ts"), - "utf8", - ); - const qaRuntimeApi = readFileSync( - path.resolve(TEST_DIR, "../../extensions/qa-lab/src/runtime-api.ts"), - "utf8", - ); - - expect(script).toContain('ln -sfnT "$openclaw_package_dir/dist" /app/dist'); - expect(script).toContain('-v "$ROOT_DIR/dist:/app/.openclaw-qa-harness-dist:ro"'); + expect(script).toContain('cp "$ROOT_DIR/package.json" "$harness_package_json"'); expect(script).toContain( - 'ln -sfnT /app/.openclaw-qa-harness-dist "$openclaw_package_dir/.openclaw-qa-harness-dist"', + 'node "$ROOT_DIR/scripts/e2e/lib/npm-telegram-live/prepare-package.mjs" "$harness_package_json"', ); - expect(script).toContain('cp "$openclaw_package_dir/package.json" /app/package.json'); - expect(script).toContain('-v "$ROOT_DIR/extensions/qa-lab:/app/extensions/qa-lab:ro"'); + expect(script).toContain('-v "$harness_package_json:/app/package.json:ro"'); + expect(script).toContain('-v "$ROOT_DIR/dist:/app/dist:ro"'); + expect(script).toContain('-v "$ROOT_DIR/node_modules:/trusted-harness/node_modules:ro"'); + expect(script).toContain('-v "$ROOT_DIR/packages:/app/packages:ro"'); + expect(script).toContain('-v "$ROOT_DIR/extensions:/app/extensions:ro"'); + expect(script).toContain('-v "$ROOT_DIR/taxonomy.yaml:/app/taxonomy.yaml:ro"'); expect(script).toContain('-v "$ROOT_DIR/qa/scenarios:/app/qa/scenarios:ro"'); - expect(script).not.toContain('ln -sfnT /app/extensions "$openclaw_package_dir/extensions"'); - expect(script).toContain("node scripts/e2e/lib/npm-telegram-live/prepare-package.mjs"); - expect(script).toContain("/app/node_modules/openclaw/package.json"); - expect(preparePackage).toContain('pkg.exports["./plugin-sdk/gateway-runtime"]'); - expect(preparePackage).toContain('"./dist/plugin-sdk/gateway-runtime.js"'); - expect(preparePackage).toContain('pkg.exports["./plugin-sdk/qa-runtime"]'); - expect(preparePackage).toContain('"./.openclaw-qa-harness-dist/plugin-sdk/qa-runtime.js"'); - expect(gatewayRpcClient).toContain('from "openclaw/plugin-sdk/gateway-runtime"'); - expect(qaRuntimeApi).toContain('from "openclaw/plugin-sdk/gateway-runtime"'); + expect(script).toContain("for dependency_dir in /trusted-harness/node_modules/*"); + expect(script).toContain("for workspace_dir in /app/packages/* /app/extensions/*"); + expect(script).toContain('link_harness_dependency "$workspace_dir" "$workspace_name"'); + expect(script).toContain("link_harness_dependency /app openclaw"); + expect(script).not.toContain('openclaw_package_dir="/npm-global/lib/node_modules/openclaw"'); + expect(script).not.toContain('cp "$openclaw_package_dir/package.json" /app/package.json'); + expect(script).not.toContain("/app/node_modules/openclaw/package.json"); + expect(script).not.toContain("link_installed_package_dependency"); }); - it("adds private harness exports only to two ephemeral manifests", () => { + it("adds private SDK exports only to the trusted harness manifest", () => { const root = mkTempRoot(); - const packageJsonPaths = ["root-package.json", "installed-package.json"].map((name) => - path.join(root, name), + const harnessManifestPath = path.join(root, "harness-package.json"); + const candidateManifestPath = path.join(root, "candidate-package.json"); + const existingGatewayExport = { + types: "./existing/gateway-runtime.d.ts", + default: "./existing/gateway-runtime.js", + }; + writeFileSync( + harnessManifestPath, + `${JSON.stringify({ + name: "openclaw", + exports: { + "./kept": "./dist/kept.js", + "./plugin-sdk/gateway-runtime": existingGatewayExport, + }, + })}\n`, ); - for (const packageJsonPath of packageJsonPaths) { - writeFileSync(packageJsonPath, JSON.stringify({ exports: { ".": "./dist/index.js" } })); - } + writeFileSync(candidateManifestPath, '{"name":"candidate","exports":{}}\n'); + const candidateBefore = readFileSync(candidateManifestPath, "utf8"); - const result = spawnSync(process.execPath, [PREPARE_PACKAGE_PATH, ...packageJsonPaths], { - encoding: "utf8", + execFileSync(process.execPath, [PREPARE_PACKAGE_PATH, harnessManifestPath]); + + const prepared = JSON.parse(readFileSync(harnessManifestPath, "utf8")) as { + exports: Record; + }; + expect(prepared.exports["./kept"]).toBe("./dist/kept.js"); + expect(prepared.exports["./plugin-sdk/gateway-runtime"]).toEqual(existingGatewayExport); + expect(prepared.exports["./plugin-sdk/qa-runtime"]).toEqual({ + default: "./dist/plugin-sdk/qa-runtime.js", }); - - expect(result.status).toBe(0); - for (const packageJsonPath of packageJsonPaths) { - const pkg = JSON.parse(readFileSync(packageJsonPath, "utf8")) as { - exports: Record; - }; - expect(pkg.exports).toMatchObject({ - ".": "./dist/index.js", - "./plugin-sdk/gateway-runtime": { - types: "./dist/plugin-sdk/gateway-runtime.d.ts", - default: "./dist/plugin-sdk/gateway-runtime.js", - }, - "./plugin-sdk/qa-runtime": { - default: "./.openclaw-qa-harness-dist/plugin-sdk/qa-runtime.js", - }, + expect(prepared.exports["./plugin-sdk/qa-lab"]).toEqual({ + default: "./dist/plugin-sdk/qa-lab.js", + }); + for (const subpath of privateLocalOnlyPluginSdkEntrypoints) { + expect(prepared.exports[`./plugin-sdk/${subpath}`]).toEqual({ + default: `./dist/plugin-sdk/${subpath}.js`, }); } - - const thirdPackageJsonPath = path.join(root, "third-package.json"); - writeFileSync(thirdPackageJsonPath, JSON.stringify({ exports: { ".": "./dist/index.js" } })); - const rejected = spawnSync( - process.execPath, - [PREPARE_PACKAGE_PATH, ...packageJsonPaths, thirdPackageJsonPath], - { encoding: "utf8" }, - ); - - expect(rejected.status).toBe(1); - expect(rejected.stderr).toContain("expected exactly two ephemeral package manifests, got 3"); - expect(JSON.parse(readFileSync(thirdPackageJsonPath, "utf8"))).toEqual({ - exports: { ".": "./dist/index.js" }, - }); - }); - - it("exposes installed package dependencies to the mounted QA harness", () => { - const script = readFileSync(DOCKER_SCRIPT_PATH, "utf8"); - - expect(script).toContain("link_installed_package_dependency()"); - expect(script).toContain( - 'local source="/npm-global/lib/node_modules/openclaw/node_modules/$name"', - ); - expect(script).toContain('ln -sfn "$source" "$target"'); - expect(script).toContain('link_installed_package_dependency "$dependency"'); - expect(script).toContain("@modelcontextprotocol/sdk"); - expect(script).toContain("yaml"); - expect(script).toContain("zod"); + expect(readFileSync(candidateManifestPath, "utf8")).toBe(candidateBefore); }); it("lets npm-specific credential aliases override shared QA env", () => { @@ -425,6 +413,41 @@ describe("package Telegram live Docker E2E", () => { }); }); + it.each([ + { + name: "promotes the default canary before taxonomy-backed release selection", + env: {}, + requested: [], + resolved: ["telegram-status-command"], + expected: ["channel-canary", "telegram-status-command"], + }, + { + name: "keeps focused non-RTT selections unchanged", + env: {}, + requested: ["telegram-status-command"], + resolved: ["telegram-status-command"], + expected: ["telegram-status-command"], + }, + { + name: "promotes an explicitly requested RTT canary", + env: { OPENCLAW_NPM_TELEGRAM_RTT_CHECKS: "channel-canary" }, + requested: ["telegram-status-command"], + resolved: ["telegram-status-command"], + expected: ["channel-canary", "telegram-status-command"], + }, + { + name: "does not duplicate an already selected RTT canary", + env: {}, + requested: ["telegram-status-command", "channel-canary"], + resolved: ["telegram-status-command", "channel-canary"], + expected: ["channel-canary", "telegram-status-command"], + }, + ])("$name", ({ env, requested, resolved, expected }) => { + const options = testing.resolveRttOptions(env, requested); + + expect(testing.prioritizeRoundTripProbeScenario(resolved, options)).toEqual(expected); + }); + it("rejects retired RTT scenario ids", () => { expect(() => testing.resolveRttOptions({