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({