feat(test): share the measured Gateway host across plugin workloads (#155830)

This commit is contained in:
Vincent Koc 2026-09-23 03:13:10 +08:00 • committed by GitHub
parent c619369ae8
commit 4bc30cabb7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 441 additions and 168 deletions

View file

@ -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.

View file

@ -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<KitchenSinkResourceCase["shutdown"]>,
) {
@ -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<unknown>;
sample: () => ReturnType<typeof readGatewayResources>;
measure: (
name: string,
count: number,
run: (index: number) => Promise<void>,
) => Promise<KitchenSinkResourcePhase>;
};
/** 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<typeof resolveResourceGatewayRuntime>;
prepare: (context: ResourceGatewayPreparation) => Promise<void>;
run: (context: ResourceGatewayWorkload) => Promise<void>;
}) {
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:<local.tgz>",
);
}
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<void>,
) => {
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");
}

View file

@ -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));
});
});

View file

@ -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);
});
},
);