From 4bc30cabb791f9af124ff038f0ce09c8c3da3d31 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Wed, 23 Sep 2026 03:13:10 +0800 Subject: [PATCH] feat(test): share the measured Gateway host across plugin workloads (#155830) --- docs/reference/test/performance.md | 24 + scripts/e2e/kitchen-sink-rpc-walk.mts | 417 +++++++++++------- .../fixtures/gateway-resource-host.mjs | 21 + test/scripts/gateway-resource-host.test.ts | 147 ++++++ 4 files changed, 441 insertions(+), 168 deletions(-) create mode 100644 test/scripts/fixtures/gateway-resource-host.mjs create mode 100644 test/scripts/gateway-resource-host.test.ts diff --git a/docs/reference/test/performance.md b/docs/reference/test/performance.md index cf376f83c342..f17aa7a6c73c 100644 --- a/docs/reference/test/performance.md +++ b/docs/reference/test/performance.md @@ -97,6 +97,30 @@ separately; post-process-exit sampling and guaranteed reclamation remain unsuppo These observations establish neither a leak nor a budget violation. Repeat comparable pairs through the campaign owner before drawing performance conclusions; do not sum individual plugin costs. +### Reusing the resource host + +Source-checkout campaign tools can import `resolveResourceGatewayRuntime` and +`runResourceGatewayCase` from `scripts/e2e/kitchen-sink-rpc-walk.mts`. Kitchen Sink +uses this same host lifecycle. Run from one frozen, built OpenClaw package root +per process; this is a testing seam, not a published plugin SDK API. + +The preparation callback receives an isolated config path, loopback port, test +token and a local-archive installer that verifies the supplied SHA-256. Enable +only the selected plugins there. The workload callback receives authenticated +CLI-mode RPC calls, resource snapshots and counted `measure(name, count, run)` +phases. Assert the active plugin inventory and operation results in the workload; +registration or a successful transport response alone does not establish coverage. + +The host records startup, preserves failed phases and joins Gateway shutdown +before checking service-stop logs. A failed workload, nonzero exit, attempted +forced cleanup or shutdown error retains temporary state and fails the case. +The campaign owns bounded callback deadlines, mock-service cleanup, runner +isolation, repetitions and report publication. Keep mock-service measurements +separate from Gateway observations and preserve host, archive and harness hashes. +On an outer timeout, the runner must terminate and join the whole container or +cgroup: the Gateway has its own process group, so killing the campaign process +alone does not clean it up. + ### Zod schema compilation Compile individual schemas only after measuring a repeated validation path. diff --git a/scripts/e2e/kitchen-sink-rpc-walk.mts b/scripts/e2e/kitchen-sink-rpc-walk.mts index 8d9acd5b02e0..1028da369cd2 100644 --- a/scripts/e2e/kitchen-sink-rpc-walk.mts +++ b/scripts/e2e/kitchen-sink-rpc-walk.mts @@ -2792,16 +2792,22 @@ export function assertKitchenSinkResourcePlugins(payload: unknown, enabled: bool return ids; } -type KitchenSinkResourceCase = { - name: "empty" | "conformance"; +export type GatewayResourceCase = { + name: string; status: "blocked" | "exercised" | "failed"; phases: KitchenSinkResourcePhase[]; activePlugins?: string[]; - calibration?: KitchenSinkCalibration; + host?: { commit: unknown; version: unknown; buildId: unknown; entrySha256: string }; + fixtures?: Array<{ archive: string; sha256: string }>; shutdown?: { exited: boolean; signals: string[]; exitCode: number | null; signal: string | null }; error?: string; }; +type KitchenSinkResourceCase = GatewayResourceCase & { + name: "empty" | "conformance"; + calibration?: KitchenSinkCalibration; +}; + export function assertKitchenSinkResourceShutdown( shutdown: NonNullable, ) { @@ -2815,7 +2821,7 @@ export function assertKitchenSinkResourceShutdown( } } -async function profileKitchenSinkResources(reportPath: string) { +export function resolveResourceGatewayRuntime() { if ( process.versions.bun || process.platform !== "linux" || @@ -2834,6 +2840,182 @@ async function profileKitchenSinkResources(reportPath: string) { "Resource comparison requires a built OpenClaw entry in the current package root", ); } + const buildInfo = asRecord(readJson(path.resolve("dist/build-info.json"))); + if (typeof buildInfo.commit !== "string" || !/^[a-f0-9]{40}$/u.test(buildInfo.commit)) { + throw new Error("Resource comparison requires dist/build-info.json with a full source commit"); + } + return { runner, buildInfo }; +} + +type ResourceGatewayPreparation = { + root: string; + env: KitchenSinkEnv; + port: number; + token: string; + installArchive: ( + file: string, + expectedSha256: string, + ) => Promise<{ archive: string; sha256: string }>; +}; + +type ResourceGatewayWorkload = ResourceGatewayPreparation & { + rpc: (method: string, params: unknown) => Promise; + sample: () => ReturnType; + measure: ( + name: string, + count: number, + run: (index: number) => Promise, + ) => Promise; +}; + +/** One frozen host package root as cwd per process; callers own outer deadlines and mock cleanup. */ +export async function runResourceGatewayCase(options: { + result: GatewayResourceCase; + runtime: ReturnType; + prepare: (context: ResourceGatewayPreparation) => Promise; + run: (context: ResourceGatewayWorkload) => Promise; +}) { + const { + result, + runtime: { runner, buildInfo }, + } = options; + delete result.error; + const { root, env }: { root: string; env: KitchenSinkEnv } = makeEnv(kitchenSinkResourceEnv()); + // The shared host must not select a fixture personality for other workloads. + delete env.OPENCLAW_KITCHEN_SINK_PERSONALITY; + const logPath = path.join(root, "gateway.log"); + let child: ChildProcess | undefined; + const fail = (error: unknown) => { + result.status = "failed"; + result.error = [result.error, error instanceof Error ? error.message : String(error)] + .filter(Boolean) + .join("; ") + .slice(0, 2_048); + }; + try { + result.host = { + commit: buildInfo.commit, + version: buildInfo.version, + buildId: buildInfo.buildId, + entrySha256: createHash("sha256").update(fs.readFileSync(runner.baseArgs[0]!)).digest("hex"), + }; + result.fixtures = []; + const port = await resolveKitchenSinkRpcPort(); + writeJson(env.OPENCLAW_CONFIG_PATH, { + gateway: { + bind: "loopback", + port, + auth: { mode: "token", token: TOKEN }, + controlUi: { enabled: false }, + }, + // A memory slot can activate a plugin outside the allowlist. + plugins: { enabled: false, slots: { memory: "none" } }, + }); + const context = { + root, + env, + port, + token: TOKEN, + installArchive: async (file: string, expectedSha256: string) => { + const archive = path.resolve(file); + if ( + !archive.endsWith(".tgz") || + !fs.statSync(archive).isFile() || + !/^[a-f0-9]{64}$/u.test(expectedSha256) || + createHash("sha256").update(fs.readFileSync(archive)).digest("hex") !== expectedSha256 + ) { + throw new Error("Resource fixture must be a local archive matching its SHA-256"); + } + const help = await runOpenClaw(runner, ["plugins", "install", "--help"], env); + if (help.stdoutTruncatedChars) { + throw new Error("Plugin fixture help probe output was truncated"); + } + await runOpenClaw( + runner, + [ + "plugins", + "install", + `npm-pack:${archive}`, + "--force", + ...fixtureCapabilityConsentArgs(help.stdout), + ], + env, + { timeoutMs: resolveKitchenSinkRpcConfig(env).installTimeoutMs }, + ); + const receipt = { archive: path.basename(archive), sha256: expectedSha256 }; + result.fixtures!.push(receipt); + return receipt; + }, + }; + await options.prepare(context); + child = await startGateway(runner, port, env, logPath, true); + await waitForGatewayReady(child, port, logPath); + const sample = () => readGatewayResources(child!); + result.phases.push( + summarizeResourcePhase( + "startup", + await readGatewayResources(child, { initial: true }), + await sample(), + { attempted: 0, completed: 0, failed: 0 }, + ), + ); + await options.run({ + ...context, + rpc: (method, params) => rpcCall(method, params, { runner, env, port }), + sample, + measure: async (name, count, run) => { + const measured = await measureResourceOperations({ name, count, sample, run }); + result.phases.push(measured); + if (measured.status !== "exercised") { + throw new Error(`${name}: ${measured.error}`); + } + return measured; + }, + }); + result.status = "exercised"; + } catch (error) { + fail(error); + } finally { + if (child) { + const signals: string[] = []; + try { + await stopGateway(child, { + killProcess: (pid, signal) => { + if (signal !== 0) { + signals.push(String(signal)); + } + return defaultKillProcess(pid, signal); + }, + }); + result.shutdown = { + exited: !isGatewayAlive(child, defaultKillProcess), + signals, + exitCode: child.exitCode, + signal: child.signalCode, + }; + assertKitchenSinkResourceShutdown(result.shutdown); + // Service-stop errors arrive during shutdown; inspect only after join. + assertNoErrorLogs(logPath); + } catch (error) { + fail(`Shutdown failed: ${String(error)}`); + } + } + if (result.status === "exercised" && process.env.OPENCLAW_KITCHEN_SINK_KEEP_TMP !== "1") { + try { + await cleanupKitchenSinkEnv(root, { throwOnFailure: true }); + } catch (error) { + fail(`Temporary state cleanup failed: ${String(error)}`); + } + } else { + console.error(`Gateway resource temp root preserved: ${root}`); + } + } + return result; +} + +async function profileKitchenSinkResources(reportPath: string) { + const runtime = resolveResourceGatewayRuntime(); + const { runner, buildInfo } = runtime; const fixture = path.resolve(PLUGIN_SPEC.slice("npm-pack:".length)); if ( !PLUGIN_SPEC.startsWith("npm-pack:") || @@ -2845,10 +3027,6 @@ async function profileKitchenSinkResources(reportPath: string) { "Resource comparison requires OPENCLAW_KITCHEN_SINK_NPM_SPEC=npm-pack:", ); } - const buildInfo = asRecord(readJson(path.resolve("dist/build-info.json"))); - if (typeof buildInfo.commit !== "string" || !/^[a-f0-9]{40}$/u.test(buildInfo.commit)) { - throw new Error("Resource comparison requires dist/build-info.json with a full source commit"); - } const sha256 = (file: string | URL) => createHash("sha256").update(fs.readFileSync(file)).digest("hex"); const fixtureSha256 = sha256(fixture); @@ -2920,176 +3098,79 @@ async function profileKitchenSinkResources(reportPath: string) { }; try { for (const result of cases) { - const { name } = result; - const enabled = name === "conformance"; - delete result.error; - const { root, env } = makeEnv(kitchenSinkResourceEnv()); - const logPath = path.join(root, "gateway.log"); - let child: ChildProcess | undefined; - try { - const port = await resolveKitchenSinkRpcPort(); - writeJson(env.OPENCLAW_CONFIG_PATH, { - gateway: { - bind: "loopback", - port, - auth: { mode: "token", token: TOKEN }, - controlUi: { enabled: false }, - }, - // The default memory slot bypasses the allowlist when plugins are enabled. - // Keep it off in both cases so conformance measures only the fixture. - plugins: { enabled: false, slots: { memory: "none" } }, - }); - if (enabled) { - if (sha256(fixture) !== fixtureSha256) { - throw new Error("Kitchen Sink fixture changed after profiling admission"); + const enabled = result.name === "conformance"; + await runResourceGatewayCase({ + result, + runtime, + prepare: async ({ env, port, installArchive }) => { + env.OPENCLAW_KITCHEN_SINK_PERSONALITY = "conformance"; + if (enabled) { + if (sha256(fixture) !== fixtureSha256) { + throw new Error("Kitchen Sink fixture changed after profiling admission"); + } + await installArchive(fixture, fixtureSha256); + configureKitchenSink(env, port); } - const help = await runOpenClaw(runner, ["plugins", "install", "--help"], env); - if (help.stdoutTruncatedChars) { - throw new Error("Plugin fixture help probe output was truncated"); - } - await runOpenClaw( - runner, - [ - "plugins", - "install", - `npm-pack:${fixture}`, - "--force", - ...fixtureCapabilityConsentArgs(help.stdout), - ], - env, - { - timeoutMs: resolveKitchenSinkRpcConfig(env).installTimeoutMs, - }, + }, + run: async ({ rpc, sample, measure }) => { + result.activePlugins = assertKitchenSinkResourcePlugins( + await rpc("plugins.list", {}), + enabled, ); - configureKitchenSink(env, port); - } - child = await startGateway(runner, port, env, logPath, true); - await waitForGatewayReady(child, port, logPath); - const sample = () => readGatewayResources(child!); - const initial = await readGatewayResources(child, { initial: true }); - result.phases.push( - summarizeResourcePhase("startup", initial, await sample(), { - attempted: 0, - completed: 0, - failed: 0, - }), - ); - const rpcOptions = { runner, env, port }; - result.activePlugins = assertKitchenSinkResourcePlugins( - await rpcCall("plugins.list", {}, rpcOptions), - enabled, - ); - // Equal warmup and neutral work isolate enabled-plugin host overhead. - // Never retry measured RPCs: a failed response must remain a failed operation. - assertGatewayHealthPayload(await rpcCall("health", {}, rpcOptions)); - const idle = async (phase: string) => { - const before = await sample(); - await delay(report.measurement.idleMs); - result.phases.push( - summarizeResourcePhase(phase, before, await sample(), { - attempted: 0, - completed: 0, - failed: 0, - }), - ); - }; - await idle("idle"); - const runPhase = async ( - phase: string, - count: number, - run: (index: number) => Promise, - ) => { - const measured = await measureResourceOperations({ name: phase, count, sample, run }); - result.phases.push(measured); - if (measured.status !== "exercised") { - throw new Error(`${phase}: ${measured.error}`); - } - }; - await runPhase("neutral-rpc", report.measurement.neutralOperations, async () => { - assertGatewayHealthPayload(await rpcCall("health", {}, rpcOptions)); - }); - await idle("post-neutral"); - if (enabled) { - const session = assertCreatedKitchenSinkSession( - await rpcCall( - "sessions.create", - { key: SESSION_KEY, agentId: "main", label: "kitchen-sink-resources" }, - rpcOptions, - ), - ); - await runPhase("plugin-tool", report.measurement.pluginToolOperations, async (index) => { - assertKitchenSinkTextInvokeResult( - await rpcCall( - "tools.invoke", - { + // Equal warmup and neutral work isolate enabled-plugin host overhead. + // Never retry measured RPCs: a failed response must remain a failed operation. + assertGatewayHealthPayload(await rpc("health", {})); + const idle = async (phase: string) => { + const before = await sample(); + await delay(report.measurement.idleMs); + result.phases.push( + summarizeResourcePhase(phase, before, await sample(), { + attempted: 0, + completed: 0, + failed: 0, + }), + ); + }; + await idle("idle"); + await measure("neutral-rpc", report.measurement.neutralOperations, async () => { + assertGatewayHealthPayload(await rpc("health", {})); + }); + await idle("post-neutral"); + if (enabled) { + const session = assertCreatedKitchenSinkSession( + await rpc("sessions.create", { + key: SESSION_KEY, + agentId: "main", + label: "kitchen-sink-resources", + }), + ); + await measure("plugin-tool", report.measurement.pluginToolOperations, async (index) => { + assertKitchenSinkTextInvokeResult( + await rpc("tools.invoke", { name: "kitchen_sink_text", args: { prompt: "explain kitchen sink resource profiling" }, sessionKey: String(session.key), agentId: "main", idempotencyKey: `kitchen-sink-resources-${index}`, - }, - rpcOptions, - ), - ); - }); - await idle("post-tool"); - result.calibration = await calibrateKitchenSinkResources({ - pluginId: PLUGIN_ID, - rpc: (method, params) => rpcCall(method, params, rpcOptions), - sample, - assertDisabled: (payload) => { - assertKitchenSinkResourcePlugins(payload, false); - }, - }); - report.postDisposalResidual = result.calibration.postDisposalResidual; - if (result.calibration.status !== "exercised") { - throw new Error(`Resource calibration: ${result.calibration.error}`); - } - } - assertNoErrorLogs(logPath); - result.status = "exercised"; - } catch (error) { - result.status = "failed"; - result.error = String(error instanceof Error ? error.message : error).slice(0, 2_048); - } finally { - if (child) { - const signals: string[] = []; - try { - await stopGateway(child, { - killProcess: (pid, signal) => { - if (signal !== 0) { - signals.push(String(signal)); - } - return defaultKillProcess(pid, signal); + }), + ); + }); + await idle("post-tool"); + result.calibration = await calibrateKitchenSinkResources({ + pluginId: PLUGIN_ID, + rpc, + sample, + assertDisabled: (payload) => { + assertKitchenSinkResourcePlugins(payload, false); }, }); - const exited = !isGatewayAlive(child, defaultKillProcess); - result.shutdown = { - exited, - signals, - exitCode: child.exitCode, - signal: child.signalCode, - }; - assertKitchenSinkResourceShutdown(result.shutdown); - } catch (error) { - result.status = "failed"; - result.error = [result.error, `Shutdown failed: ${String(error)}`] - .filter(Boolean) - .join("; ") - .slice(0, 2_048); + report.postDisposalResidual = result.calibration.postDisposalResidual; + if (result.calibration.status !== "exercised") { + throw new Error(`Resource calibration: ${result.calibration.error}`); + } } - } - if (result.status === "exercised" && process.env.OPENCLAW_KITCHEN_SINK_KEEP_TMP !== "1") { - try { - await cleanupKitchenSinkEnv(root, { throwOnFailure: true }); - } catch (error) { - result.status = "failed"; - result.error = `Temporary state cleanup failed: ${String(error)}`.slice(0, 2_048); - } - } else { - console.error(`Kitchen Sink resource temp root preserved: ${root}`); - } - } + }, + }); if (result.status !== "exercised") { throw new Error(result.error ?? "Kitchen Sink resource case did not complete"); } diff --git a/test/scripts/fixtures/gateway-resource-host.mjs b/test/scripts/fixtures/gateway-resource-host.mjs new file mode 100644 index 000000000000..65a14aa396fe --- /dev/null +++ b/test/scripts/fixtures/gateway-resource-host.mjs @@ -0,0 +1,21 @@ +import { appendFileSync } from "node:fs"; +import { createServer } from "node:http"; + +if (process.env.FIXTURE_INVOCATIONS) { + appendFileSync(process.env.FIXTURE_INVOCATIONS, JSON.stringify(process.argv) + "\n"); +} +const port = Number(process.argv[process.argv.indexOf("--port") + 1]); +const server = createServer((request, response) => { + response.setHeader("content-type", "application/json"); + response.end(JSON.stringify(request.url === "/readyz" ? { ready: true } : { completed: true })); +}); +server.listen(port, "127.0.0.1"); +process.once("SIGTERM", () => { + if (process.env.FIXTURE_STOP_ERROR === "1") { + console.error("[error] fixture service stop failed"); + } + server.close(() => { + process.disconnect(); + process.exit(Number(process.env.FIXTURE_EXIT_CODE ?? 0)); + }); +}); diff --git a/test/scripts/gateway-resource-host.test.ts b/test/scripts/gateway-resource-host.test.ts new file mode 100644 index 000000000000..91afb81851e3 --- /dev/null +++ b/test/scripts/gateway-resource-host.test.ts @@ -0,0 +1,147 @@ +import { existsSync, rmSync, writeFileSync } from "node:fs"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { afterEach, describe, expect, it } from "vitest"; +import { + runResourceGatewayCase, + type GatewayResourceCase, +} from "../../scripts/e2e/kitchen-sink-rpc-walk.mts"; + +const roots: string[] = []; +afterEach(() => { + for (const root of roots.splice(0)) { + rmSync(root, { recursive: true, force: true }); + } +}); + +const runtime = { + runner: { + command: process.execPath, + baseArgs: [fileURLToPath(new URL("./fixtures/gateway-resource-host.mjs", import.meta.url))], + }, + buildInfo: { commit: "a".repeat(40), version: "fixture", buildId: "fixture-build" }, +}; + +function result(): GatewayResourceCase { + return { name: "fixture", status: "blocked", phases: [] }; +} + +describe.skipIf(process.platform !== "linux" || typeof process.threadCpuUsage !== "function")( + "shared resource Gateway lifecycle", + () => { + it("records actual child samples and joins a successful workload before removing state", async () => { + const receipt = result(); + let root = ""; + await runResourceGatewayCase({ + result: receipt, + runtime, + prepare: async (context) => { + root = context.root; + roots.push(root); + expect(context.env.OPENCLAW_KITCHEN_SINK_PERSONALITY).toBeUndefined(); + }, + run: async ({ port, measure }) => { + const url = new URL("http://127.0.0.1/work"); + url.port = String(port); + await measure("work", 2, async () => { + const response = await fetch(url); + expect(await response.json()).toEqual({ completed: true }); + }); + }, + }); + expect(receipt.status).toBe("exercised"); + expect(receipt.phases.map(({ name }) => name)).toEqual(["startup", "work"]); + expect(receipt.phases[1]?.operations).toEqual({ attempted: 2, completed: 2, failed: 0 }); + expect(receipt.phases[1]?.before.pid).toBe(receipt.phases[1]?.after?.pid); + expect(receipt.shutdown).toMatchObject({ exited: true, exitCode: 0, signal: null }); + expect(receipt.host).toMatchObject(runtime.buildInfo); + expect(existsSync(root)).toBe(false); + }); + + it("retains preparation failure without starting a child or running work", async () => { + const receipt = result(); + await runResourceGatewayCase({ + result: receipt, + runtime, + prepare: async ({ root }) => { + roots.push(root); + throw new Error("fixture preparation failed"); + }, + run: async () => { + throw new Error("work must not run"); + }, + }); + expect(receipt).toMatchObject({ + status: "failed", + error: "fixture preparation failed", + phases: [], + }); + expect(receipt.shutdown).toBeUndefined(); + expect(existsSync(roots[0]!)).toBe(true); + }); + + it.each([ + { workloadFails: true, exitCode: "0", stopError: "0" }, + { workloadFails: false, exitCode: "1", stopError: "0" }, + { workloadFails: false, exitCode: "0", stopError: "1" }, + { workloadFails: true, exitCode: "1", stopError: "0" }, + ])( + "preserves work and invalidates failure $workloadFails/$exitCode/$stopError", + async ({ workloadFails, exitCode, stopError }) => { + const receipt = result(); + await runResourceGatewayCase({ + result: receipt, + runtime, + prepare: async ({ root, env }) => { + roots.push(root); + env.FIXTURE_EXIT_CODE = exitCode; + env.FIXTURE_STOP_ERROR = stopError; + }, + run: async ({ measure }) => { + await measure("completed", 1, async () => {}); + if (workloadFails) { + await measure("failed", 1, async () => { + throw new Error("fixture assertion failed"); + }); + } + }, + }); + expect(receipt.status).toBe("failed"); + expect(receipt.phases[1]?.operations.completed).toBe(1); + expect(receipt.shutdown).toMatchObject({ exited: true, exitCode: Number(exitCode) }); + if (workloadFails) { + expect(receipt.error).toContain("fixture assertion failed"); + expect(receipt.phases[2]?.operations).toEqual({ attempted: 1, completed: 0, failed: 1 }); + } + if (exitCode === "1") { + expect(receipt.error).toContain("did not exit cleanly"); + } + if (stopError === "1") { + expect(receipt.error).toContain("fixture service stop failed"); + } + expect(existsSync(roots[0]!)).toBe(true); + }, + ); + + it("rejects a mismatched archive before invoking the installer", async () => { + const receipt = result(); + let invocations = ""; + await runResourceGatewayCase({ + result: receipt, + runtime, + prepare: async ({ root, env, installArchive }) => { + roots.push(root); + const archive = path.join(root, "fixture.tgz"); + invocations = path.join(root, "invocations"); + env.FIXTURE_INVOCATIONS = invocations; + writeFileSync(archive, "mismatched fixture"); + await installArchive(archive, "0".repeat(64)); + }, + run: async () => {}, + }); + expect(receipt).toMatchObject({ status: "failed", fixtures: [], phases: [] }); + expect(receipt.error).toContain("matching its SHA-256"); + expect(existsSync(invocations)).toBe(false); + }); + }, +);