mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
test(gateway): add sustained heap retention diagnostics (#163764)
* test(gateway): add sustained heap retention diagnostics Add a synthetic dist Gateway workload with signed webchat clients, mock inference, Code Mode, verified subagent completion, minute memory samples, and paired native heap snapshots. Extend the streaming snapshot comparator with named strong paths and separate dominator chains. Keep retryable project-access refusals visible and revalidate turn admission after asynchronous preparation. The new admission regression fails against the exact prior helper and passes after the repair. Remote checks and independent review passed. Measurements do not establish a production leak fix. * fix(tooling): register Gateway heap rig entry point
This commit is contained in:
parent
12e5f8dac6
commit
05ce3c3a1a
9 changed files with 1207 additions and 28 deletions
|
|
@ -176,6 +176,8 @@ const repositoryScriptEntries = [
|
|||
"scripts/fixtures/packed-plugin-sdk-type-smoke.ts!",
|
||||
// Generates the native browser page scripts from their UI source modules.
|
||||
"scripts/generate-browser-inspect-script-swift.mts!",
|
||||
// The diagnostics guide invokes the sustained Gateway heap rig by path.
|
||||
"scripts/gateway-heap-rig.mjs!",
|
||||
// The diagnostics guide invokes this offline snapshot comparison CLI by path.
|
||||
"scripts/heap-snapshot-diff.mjs!",
|
||||
// CI executes screenshot evidence from the workflow-owned harness copy.
|
||||
|
|
|
|||
|
|
@ -345,6 +345,33 @@ paths and V8-specific weak/ephemeron semantics; the script is a strong-edge grap
|
|||
summary. `--json` produces machine-readable output. Treat diff output as sensitive
|
||||
too: it contains unredacted heap names.
|
||||
|
||||
Add `--top 40 --max-depth 60` to include named strong retaining paths and
|
||||
dominator chains for the largest growers. `--node <id>` selects a particular
|
||||
object in the later snapshot. A shortest root path shows reachability;
|
||||
the separate dominator chain identifies exclusive retention in that graph.
|
||||
|
||||
From a built source checkout, an isolated synthetic workload can collect a
|
||||
comparable pair without connecting to an existing Gateway:
|
||||
|
||||
```bash
|
||||
node scripts/gateway-heap-rig.mjs --root .rig/node26 --minutes 90
|
||||
```
|
||||
|
||||
Run this on a dedicated host with enough memory for snapshots. It starts the
|
||||
dist Gateway and local mock model servers on loopback ports 19548–19550,
|
||||
seeds 2,000 sessions across two agents, and drives ten reconnecting Control UI
|
||||
WebSocket clients plus mock model, Code Mode, and subagent turns. A synthetic
|
||||
catalog plugin exercises Gateway projection and publication ownership; it does
|
||||
not emulate a native provider's caches or remote-node transport.
|
||||
The root must be new. All state, logs, minute samples, and snapshots stay there.
|
||||
The rig stops its children on completion or interruption and retains evidence.
|
||||
Successful RPC counts and any retried refusals are recorded separately. If
|
||||
`projects.list` refuses a read because access facts changed, the rig retries it
|
||||
once; a second refusal or another error stops the run.
|
||||
Use the same script and settings with another Node binary for a runtime control;
|
||||
choose another root and three-port block for each run. Raw minute samples include
|
||||
allocation churn; compare the snapshot `heapUsedAfter` anchors for post-GC growth.
|
||||
|
||||
## Sampling heap profile
|
||||
|
||||
An operator with `operator.admin` can sample allocations in the Gateway's main
|
||||
|
|
|
|||
535
scripts/gateway-heap-rig.mjs
Normal file
535
scripts/gateway-heap-rig.mjs
Normal file
|
|
@ -0,0 +1,535 @@
|
|||
#!/usr/bin/env node
|
||||
// Long-running, synthetic-only Gateway retention probe. Run on an isolated host.
|
||||
import { spawn } from "node:child_process";
|
||||
import { createHash, generateKeyPairSync, randomUUID, sign } from "node:crypto";
|
||||
import {
|
||||
appendFileSync,
|
||||
closeSync,
|
||||
mkdirSync,
|
||||
openSync,
|
||||
readFileSync,
|
||||
writeFileSync,
|
||||
} from "node:fs";
|
||||
import path from "node:path";
|
||||
import { performance } from "node:perf_hooks";
|
||||
import { setTimeout as delay } from "node:timers/promises";
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { parseArgs } from "node:util";
|
||||
import { buildDeviceAuthPayloadV3 } from "../packages/gateway-client/src/device-auth.ts";
|
||||
import { applyMockOpenAiModelConfig } from "./e2e/lib/fixtures/mock-openai-config.mjs";
|
||||
import { stopChild } from "./lib/gateway-bench-child.ts";
|
||||
import { readGatewayMemory } from "./lib/gateway-bench-probes.ts";
|
||||
import { BASE_GATEWAY_BENCH_CONFIG, createGatewayBenchEnv } from "./lib/gateway-bench-runtime.ts";
|
||||
import { configureHeapRigCatalog } from "./lib/gateway-heap-rig-catalog.mjs";
|
||||
import { configureHeapRigTurns } from "./lib/gateway-heap-rig-turns.mjs";
|
||||
import { createGatewayWsClient } from "./lib/gateway-ws-client.ts";
|
||||
|
||||
const { values } = parseArgs({
|
||||
options: {
|
||||
root: { type: "string" },
|
||||
minutes: { type: "string", default: "90" },
|
||||
sessions: { type: "string", default: "2000" },
|
||||
clients: { type: "string", default: "10" },
|
||||
port: { type: "string", default: "19548" },
|
||||
"rpc-interval-ms": { type: "string", default: "2000" },
|
||||
"turn-interval-ms": { type: "string", default: "180000" },
|
||||
"snapshot-minutes": { type: "string" },
|
||||
help: { type: "boolean" },
|
||||
},
|
||||
});
|
||||
if (values.help) {
|
||||
console.log(
|
||||
"node scripts/gateway-heap-rig.mjs --root .rig/node26 --minutes 90 [--sessions 2000 --clients 10 --port 19548 --snapshot-minutes 10,90]",
|
||||
);
|
||||
process.exit(0);
|
||||
}
|
||||
function positiveInteger(value, label) {
|
||||
if (!/^[1-9]\d*$/u.test(value ?? "") || !Number.isSafeInteger(Number(value))) {
|
||||
throw new Error(`${label} must be a positive integer`);
|
||||
}
|
||||
return Number(value);
|
||||
}
|
||||
const minutes = positiveInteger(values.minutes, "minutes");
|
||||
const sessionCount = positiveInteger(values.sessions, "sessions");
|
||||
const clientCount = positiveInteger(values.clients, "clients");
|
||||
const port = positiveInteger(values.port, "port");
|
||||
const rpcInterval = positiveInteger(values["rpc-interval-ms"], "rpc-interval-ms");
|
||||
const turnInterval = positiveInteger(values["turn-interval-ms"], "turn-interval-ms");
|
||||
if (port < 19500 || port > 19597) {
|
||||
throw new Error("port and two adjacent mock ports must fit 19500–19599");
|
||||
}
|
||||
if (!values.root) {
|
||||
throw new Error("--root is required; it must name a new synthetic state directory");
|
||||
}
|
||||
const root = path.resolve(values.root);
|
||||
const build = JSON.parse(readFileSync("dist/build-info.json", "utf8"));
|
||||
const snapshotMinutes = (values["snapshot-minutes"] ?? `10,${minutes}`)
|
||||
.split(",")
|
||||
.map((v) => positiveInteger(v, "snapshot-minutes"));
|
||||
if (snapshotMinutes.some((v, i) => v > minutes || (i > 0 && v - snapshotMinutes[i - 1] < 2))) {
|
||||
throw new Error("snapshot minutes must increase by at least two and fit the duration");
|
||||
}
|
||||
mkdirSync(path.dirname(root), { recursive: true });
|
||||
mkdirSync(root); // Never reuse an operator directory or a previous run's databases.
|
||||
const eventsPath = path.join(root, "measurements.jsonl");
|
||||
const record = (event) => {
|
||||
const line = JSON.stringify({ time: new Date().toISOString(), ...event });
|
||||
appendFileSync(eventsPath, `${line}\n`);
|
||||
console.log(line);
|
||||
};
|
||||
const config = structuredClone(BASE_GATEWAY_BENCH_CONFIG);
|
||||
config.cron = { enabled: false };
|
||||
config.gateway.controlUi.allowedOrigins = [`http://127.0.0.1:${port}`];
|
||||
config.memory = { search: { enabled: false } };
|
||||
config.plugins.slots = { memory: "none" };
|
||||
applyMockOpenAiModelConfig(config, { mockPort: port + 1 });
|
||||
// Utility summaries must not consume the agent mock's ordered tool-response script.
|
||||
config.models.providers["heap-rig-utility"] = {
|
||||
...config.models.providers.openai,
|
||||
baseUrl: `http://127.0.0.1:${port + 2}/v1`,
|
||||
};
|
||||
config.agents.defaults.utilityModel = "heap-rig-utility/gpt-5.6-luna";
|
||||
config.agents.defaults = {
|
||||
...config.agents.defaults,
|
||||
heartbeat: { every: "0m" },
|
||||
skipBootstrap: true,
|
||||
skills: [],
|
||||
modelPolicy: {},
|
||||
systemAgent: { agentId: "main" },
|
||||
};
|
||||
config.agents.ownership = "explicit";
|
||||
config.agents.entries = Object.fromEntries(
|
||||
["main", "second"].map((id) => {
|
||||
const workspace = path.join(root, `workspace-${id}`);
|
||||
mkdirSync(workspace);
|
||||
return [id, { workspace }];
|
||||
}),
|
||||
);
|
||||
const turnDriver = await configureHeapRigTurns(config, root, port + 1);
|
||||
await configureHeapRigCatalog(config, root);
|
||||
const configPath = path.join(root, "openclaw.json");
|
||||
writeFileSync(configPath, `${JSON.stringify(config, null, 2)}\n`);
|
||||
const env = createGatewayBenchEnv(root, configPath, {
|
||||
startupTrace: false,
|
||||
caseEnv: {
|
||||
...turnDriver.env,
|
||||
OPENAI_API_KEY: "synthetic-heap-rig-not-a-real-key",
|
||||
PATH: `${path.dirname(process.execPath)}:${process.env.PATH}`,
|
||||
TMPDIR: root,
|
||||
},
|
||||
});
|
||||
const children = [];
|
||||
const clients = new Set();
|
||||
const requests = {};
|
||||
const retries = {};
|
||||
const events = {};
|
||||
const identities = Array.from({ length: clientCount }, () => {
|
||||
const { publicKey, privateKey } = generateKeyPairSync("ed25519");
|
||||
const rawKey = publicKey.export({ format: "jwk" }).x;
|
||||
return {
|
||||
privateKey,
|
||||
publicKey: rawKey,
|
||||
instanceId: randomUUID(),
|
||||
deviceId: createHash("sha256").update(Buffer.from(rawKey, "base64url")).digest("hex"),
|
||||
};
|
||||
});
|
||||
const inFlight = new Set();
|
||||
const unpausedWork = Promise.resolve();
|
||||
let workloadReady = unpausedWork;
|
||||
/** @type {Error | undefined} */
|
||||
let failure;
|
||||
const abort = new AbortController();
|
||||
function fail(error) {
|
||||
failure ??= error instanceof Error ? error : new Error(String(error));
|
||||
abort.abort(failure);
|
||||
}
|
||||
const onSignal = () => fail(new Error("rig interrupted"));
|
||||
process.once("SIGINT", onSignal);
|
||||
process.once("SIGTERM", onSignal);
|
||||
function start(args, label, extraEnv = {}) {
|
||||
const fd = openSync(path.join(root, `${label}.log`), "wx");
|
||||
const child = spawn(process.execPath, args, {
|
||||
cwd: process.cwd(),
|
||||
env: { ...env, ...extraEnv },
|
||||
detached: true,
|
||||
stdio: ["ignore", fd, fd],
|
||||
});
|
||||
closeSync(fd);
|
||||
children.push(child);
|
||||
child.once("error", fail);
|
||||
child.once("exit", (code, signal) => {
|
||||
if (!abort.signal.aborted) {
|
||||
fail(new Error(`${label} exited: ${code ?? signal}`));
|
||||
}
|
||||
});
|
||||
record({ kind: "process", label, pid: child.pid, node: process.version, args });
|
||||
return child;
|
||||
}
|
||||
async function sleep(ms) {
|
||||
try {
|
||||
await delay(Math.max(0, ms), undefined, { signal: abort.signal });
|
||||
} catch (error) {
|
||||
if (!abort.signal.aborted) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
class GatewayRpcError extends Error {
|
||||
constructor(method, responseError) {
|
||||
super(`${method}: ${JSON.stringify(responseError)}`);
|
||||
this.projectAccessChanged =
|
||||
method === "projects.list" &&
|
||||
responseError?.code === "UNAVAILABLE" &&
|
||||
responseError?.message ===
|
||||
"Project access changed while preparing the listing. Retry the request.";
|
||||
}
|
||||
}
|
||||
async function connect(protocol, identity) {
|
||||
const controlUi = Boolean(identity);
|
||||
const challenge = Promise.withResolvers();
|
||||
const client = createGatewayWsClient({
|
||||
url: `ws://127.0.0.1:${port}`,
|
||||
...(controlUi ? { origin: `http://127.0.0.1:${port}` } : {}),
|
||||
onEvent: ({ event, payload }) => {
|
||||
events[event] = (events[event] ?? 0) + 1;
|
||||
if (event === "connect.challenge") {
|
||||
challenge.resolve(payload);
|
||||
}
|
||||
},
|
||||
});
|
||||
clients.add(client);
|
||||
const rpc = async (method, params, timeout = 180000) => {
|
||||
const work = client.request(method, params, timeout);
|
||||
inFlight.add(work);
|
||||
try {
|
||||
const result = await work;
|
||||
if (!result.ok) {
|
||||
throw new GatewayRpcError(method, result.error);
|
||||
}
|
||||
requests[method] = (requests[method] ?? 0) + 1;
|
||||
return result.payload;
|
||||
} finally {
|
||||
inFlight.delete(work);
|
||||
}
|
||||
};
|
||||
const close = () => {
|
||||
client.close();
|
||||
clients.delete(client);
|
||||
};
|
||||
try {
|
||||
await client.waitOpen();
|
||||
const scopes = ["operator.read", "operator.write", "operator.admin"];
|
||||
if (controlUi) {
|
||||
scopes.push("operator.approvals", "operator.questions", "operator.pairing");
|
||||
}
|
||||
const params = {
|
||||
minProtocol: protocol,
|
||||
maxProtocol: protocol,
|
||||
client: {
|
||||
id: controlUi ? "openclaw-control-ui" : "gateway-client",
|
||||
displayName: "synthetic-heap-rig",
|
||||
version: "1",
|
||||
...(controlUi
|
||||
? { buildId: build.buildId, instanceId: identity.instanceId, deviceFamily: "desktop" }
|
||||
: {}),
|
||||
platform: controlUi ? "web" : process.platform,
|
||||
mode: controlUi ? "webchat" : "backend",
|
||||
},
|
||||
role: "operator",
|
||||
scopes,
|
||||
caps: controlUi ? ["tool-events", "chat-only-assistant-text", "model-selection-policy"] : [],
|
||||
};
|
||||
if (identity) {
|
||||
const timer = setTimeout(
|
||||
() => challenge.reject(new Error("Missing connect challenge")),
|
||||
10000,
|
||||
);
|
||||
let nonce;
|
||||
try {
|
||||
({ nonce } = await challenge.promise);
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
const signedAt = Date.now();
|
||||
const payload = buildDeviceAuthPayloadV3({
|
||||
deviceId: identity.deviceId,
|
||||
clientId: params.client.id,
|
||||
clientMode: params.client.mode,
|
||||
role: params.role,
|
||||
scopes,
|
||||
signedAtMs: signedAt,
|
||||
nonce,
|
||||
platform: params.client.platform,
|
||||
deviceFamily: params.client.deviceFamily,
|
||||
});
|
||||
params.device = {
|
||||
id: identity.deviceId,
|
||||
publicKey: identity.publicKey,
|
||||
signedAt,
|
||||
nonce,
|
||||
signature: sign(null, Buffer.from(payload), identity.privateKey).toString("base64url"),
|
||||
};
|
||||
}
|
||||
await rpc("connect", params);
|
||||
await rpc("sessions.subscribe", {});
|
||||
return { rpc, close };
|
||||
} catch (error) {
|
||||
close();
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
const sessions = Array.from({ length: sessionCount }, (_, index) => {
|
||||
const agentId = index % 2 === 0 ? "main" : "second";
|
||||
return { agentId, key: `agent:${agentId}:heap-rig-${index}` };
|
||||
});
|
||||
let admin;
|
||||
const tasks = [];
|
||||
try {
|
||||
const { PROTOCOL_VERSION } = await import(
|
||||
pathToFileURL(path.resolve("dist/gateway/protocol/index.js")).href
|
||||
);
|
||||
start(["scripts/e2e/mock-openai-server.mjs"], "mock", {
|
||||
MOCK_PORT: String(port + 1),
|
||||
MOCK_BIND_HOST: "127.0.0.1",
|
||||
});
|
||||
start(["scripts/e2e/mock-openai-server.mjs"], "mock-utility", {
|
||||
MOCK_PORT: String(port + 2),
|
||||
MOCK_BIND_HOST: "127.0.0.1",
|
||||
MOCK_RESPONSE_CONTROL: "",
|
||||
MOCK_REQUEST_LOG: path.join(root, "mock-utility-requests.jsonl"),
|
||||
SUCCESS_MARKER: "Synthetic activity summary",
|
||||
});
|
||||
const gateway = start(["dist/index.js", "gateway", "--port", String(port)], "gateway");
|
||||
const startupDeadline = Date.now() + 600000;
|
||||
while (!admin && !abort.signal.aborted && Date.now() < startupDeadline) {
|
||||
try {
|
||||
admin = await connect(PROTOCOL_VERSION);
|
||||
} catch {
|
||||
await sleep(1000);
|
||||
}
|
||||
}
|
||||
if (!admin) {
|
||||
throw failure ?? new Error("Gateway readiness deadline exceeded; see gateway.log");
|
||||
}
|
||||
record({
|
||||
kind: "ready",
|
||||
pid: gateway.pid,
|
||||
node: process.version,
|
||||
v8: process.versions.v8,
|
||||
build,
|
||||
sessionCount,
|
||||
clientCount,
|
||||
});
|
||||
let nextSeed = 0;
|
||||
await Promise.all(
|
||||
Array.from({ length: Math.min(8, sessionCount) }, async () => {
|
||||
for (;;) {
|
||||
const index = nextSeed++;
|
||||
if (index >= sessions.length || abort.signal.aborted) {
|
||||
return;
|
||||
}
|
||||
const { key, agentId } = sessions[index];
|
||||
await admin.rpc("sessions.create", { key, agentId });
|
||||
await admin.rpc("chat.inject", {
|
||||
sessionKey: key,
|
||||
message: `Synthetic history ${index}. `.repeat(32),
|
||||
});
|
||||
if ((index + 1) % 250 === 0) {
|
||||
record({ kind: "seed", completed: index + 1 });
|
||||
}
|
||||
}
|
||||
}),
|
||||
);
|
||||
const inventory = await admin.rpc("sessions.list", { limit: 1, archived: "all" });
|
||||
if (inventory.totalCount < sessionCount) {
|
||||
throw new Error(`Seed inventory incomplete: ${inventory.totalCount}`);
|
||||
}
|
||||
record({ kind: "seeded", totalCount: inventory.totalCount });
|
||||
const startedAt = performance.now();
|
||||
const deadline = startedAt + minutes * 60000;
|
||||
const samples = [];
|
||||
const snapshots = [];
|
||||
/** @type {Promise<unknown>} */
|
||||
let currentTurn = Promise.resolve();
|
||||
async function sample() {
|
||||
const observation = await readGatewayMemory(admin.rpc, startedAt);
|
||||
samples.push(observation);
|
||||
record({ kind: "sample", ...observation, requests: { ...requests }, retries: { ...retries } });
|
||||
}
|
||||
async function snapshot(minute) {
|
||||
/** @type {PromiseWithResolvers<void>} */
|
||||
const pause = Promise.withResolvers();
|
||||
workloadReady = pause.promise;
|
||||
try {
|
||||
await currentTurn;
|
||||
await Promise.allSettled(inFlight);
|
||||
const result = await admin.rpc(
|
||||
"diagnostics.heapSnapshot",
|
||||
{ reason: `synthetic retention rig minute ${minute}` },
|
||||
300000,
|
||||
);
|
||||
const entry = { kind: "snapshot", minute, atMs: performance.now() - startedAt, ...result };
|
||||
snapshots.push(entry);
|
||||
record(entry);
|
||||
} finally {
|
||||
workloadReady = unpausedWork;
|
||||
pause.resolve();
|
||||
}
|
||||
}
|
||||
for (let index = 0; index < clientCount; index++) {
|
||||
tasks.push(
|
||||
(async () => {
|
||||
let client;
|
||||
let reconnectAt = 0;
|
||||
let round = 0;
|
||||
try {
|
||||
while (!abort.signal.aborted && performance.now() < deadline) {
|
||||
await workloadReady;
|
||||
if (abort.signal.aborted || performance.now() >= deadline) {
|
||||
break;
|
||||
}
|
||||
if (!client || performance.now() >= reconnectAt) {
|
||||
client?.close();
|
||||
client = await connect(PROTOCOL_VERSION, identities[index]);
|
||||
reconnectAt = performance.now() + 180000 + index * 11000;
|
||||
}
|
||||
const selected =
|
||||
sessions[(Math.floor(round / 11) * clientCount + index) % sessions.length];
|
||||
const { key, agentId } = selected;
|
||||
const identity = { sessionKey: key, agentId };
|
||||
const calls = [
|
||||
[
|
||||
"sessions.list",
|
||||
{
|
||||
limit: 100,
|
||||
offset: (Math.floor(round / 11) * 100) % sessionCount,
|
||||
archived: "all",
|
||||
includeDerivedTitles: true,
|
||||
includeLastMessage: true,
|
||||
includeActivitySummary: true,
|
||||
},
|
||||
],
|
||||
[
|
||||
"sessions.catalog.list",
|
||||
{ agentId, limitPerHost: 100, progressId: randomUUID(), allowPartialResults: true },
|
||||
],
|
||||
["chat.history", { ...identity, limit: 100, maxBytes: 200000 }],
|
||||
["chat.metadata", identity],
|
||||
["models.list", { ...identity, includeDetails: true }],
|
||||
["sessions.messages.subscribe", { key, agentId, subscriptionId: "rig-active" }],
|
||||
["agent.identity.get", identity],
|
||||
["cron.list", { includeDisabled: true, limit: 50, compact: true }],
|
||||
["cron.status", {}],
|
||||
["projects.list", { includeObserved: true }],
|
||||
["sessions.messages.unsubscribe", { key, agentId, subscriptionId: "rig-active" }],
|
||||
];
|
||||
const [method, params] = calls[round++ % calls.length];
|
||||
let result;
|
||||
try {
|
||||
result = await client.rpc(method, params);
|
||||
} catch (error) {
|
||||
if (!(error instanceof GatewayRpcError) || !error.projectAccessChanged) {
|
||||
throw error;
|
||||
}
|
||||
// The projects owner refuses a read when access facts change across its await.
|
||||
retries[method] = (retries[method] ?? 0) + 1;
|
||||
record({
|
||||
kind: "rpc-retry",
|
||||
method,
|
||||
atMs: performance.now() - startedAt,
|
||||
reason: error.message,
|
||||
});
|
||||
await sleep(rpcInterval);
|
||||
await workloadReady;
|
||||
if (abort.signal.aborted || performance.now() >= deadline) {
|
||||
break;
|
||||
}
|
||||
result = await client.rpc(method, params);
|
||||
}
|
||||
if (method === "sessions.catalog.list") {
|
||||
const fixture = result.catalogs?.find((catalog) => catalog.id === "heap-rig-catalog");
|
||||
if (!fixture?.hosts?.some((host) => host.sessions.length > 0)) {
|
||||
throw new Error(
|
||||
`Synthetic catalog returned no sessions: ${JSON.stringify(result)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
await sleep(rpcInterval);
|
||||
}
|
||||
} finally {
|
||||
client?.close();
|
||||
}
|
||||
})().catch(fail),
|
||||
);
|
||||
}
|
||||
tasks.push(
|
||||
(async () => {
|
||||
for (let index = 0; !abort.signal.aborted && performance.now() < deadline; index++) {
|
||||
await workloadReady;
|
||||
if (abort.signal.aborted || performance.now() >= deadline) {
|
||||
break;
|
||||
}
|
||||
const turnWork = turnDriver.runTurn(admin.rpc, index, {
|
||||
agentId: index % 2 === 0 ? "main" : "second",
|
||||
canStart: () => !abort.signal.aborted && performance.now() < deadline,
|
||||
});
|
||||
currentTurn = turnWork;
|
||||
const turn = await turnWork;
|
||||
if (turn === null) {
|
||||
record({ kind: "turn-not-started", index, reason: "workload-ended" });
|
||||
break;
|
||||
}
|
||||
record({ ...turn, kind: "turn", turnKind: turn.kind, atMs: performance.now() - startedAt });
|
||||
await sleep(Math.min(turnInterval, deadline - performance.now()));
|
||||
}
|
||||
})().catch(fail),
|
||||
);
|
||||
await sample();
|
||||
for (let minute = 1; minute <= minutes && !abort.signal.aborted; minute++) {
|
||||
await sleep(startedAt + minute * 60000 - performance.now());
|
||||
if (abort.signal.aborted) {
|
||||
break;
|
||||
}
|
||||
if (snapshotMinutes.includes(minute)) {
|
||||
await snapshot(minute);
|
||||
}
|
||||
await sample();
|
||||
}
|
||||
if (failure) {
|
||||
throw failure;
|
||||
}
|
||||
record({
|
||||
kind: "complete",
|
||||
node: process.version,
|
||||
requests,
|
||||
retries,
|
||||
events,
|
||||
snapshots: snapshots.map(({ path: file, heapUsedAfter, atMs }) => ({
|
||||
path: file,
|
||||
heapUsedAfter,
|
||||
atMs,
|
||||
})),
|
||||
});
|
||||
writeFileSync(
|
||||
path.join(root, "result.json"),
|
||||
`${JSON.stringify({ node: process.version, v8: process.versions.v8, build, pid: gateway.pid, minutes, sessionCount, clientCount, requests, retries, events, samples, snapshots }, null, 2)}\n`,
|
||||
);
|
||||
} catch (error) {
|
||||
record({ kind: "failed", error: String(error) });
|
||||
process.exitCode = 1;
|
||||
} finally {
|
||||
abort.abort();
|
||||
for (const client of clients) {
|
||||
client.close();
|
||||
}
|
||||
await Promise.allSettled(tasks);
|
||||
for (const child of children.toReversed()) {
|
||||
record({
|
||||
kind: "cleanup",
|
||||
pid: child.pid,
|
||||
...(await stopChild(child, { teardownGraceMs: 30000 })),
|
||||
});
|
||||
}
|
||||
process.removeListener("SIGINT", onSignal);
|
||||
process.removeListener("SIGTERM", onSignal);
|
||||
}
|
||||
|
|
@ -93,7 +93,7 @@ function readGraph(file) {
|
|||
const offsets = ["type", "name", "id", "self_size", "edge_count"].map((key) =>
|
||||
fields.indexOf(key),
|
||||
);
|
||||
const edgeOffsets = ["type", "to_node"].map((key) => edgeFields.indexOf(key));
|
||||
const edgeOffsets = ["type", "to_node", "name_or_index"].map((key) => edgeFields.indexOf(key));
|
||||
if (
|
||||
offsets.includes(-1) ||
|
||||
edgeOffsets.includes(-1) ||
|
||||
|
|
@ -128,6 +128,7 @@ function readGraph(file) {
|
|||
const weak = meta.edge_types[edgeOffsets[0]].indexOf("weak");
|
||||
const shortcut = meta.edge_types[edgeOffsets[0]].indexOf("shortcut");
|
||||
let edgeType = 0;
|
||||
let edgeTarget = 0;
|
||||
reader.find('"edges"');
|
||||
const edges = reader.array((value, index) => {
|
||||
const field = index % edgeFields.length;
|
||||
|
|
@ -138,9 +139,12 @@ function readGraph(file) {
|
|||
if (value % fields.length || value / fields.length >= count) {
|
||||
throw new Error("Invalid edge target");
|
||||
}
|
||||
edgeTarget = value / fields.length;
|
||||
}
|
||||
if (field === edgeFields.length - 1) {
|
||||
// Shortcut edges are synthetic debugger conveniences, not additional retainers.
|
||||
targets[Math.floor(index / edgeFields.length)] =
|
||||
edgeType === weak || edgeType === shortcut ? count : value / fields.length;
|
||||
edgeType === weak || edgeType === shortcut ? count : edgeTarget;
|
||||
}
|
||||
});
|
||||
if (edges !== edgeCount * edgeFields.length) {
|
||||
|
|
@ -168,14 +172,14 @@ function readGraph(file) {
|
|||
}
|
||||
names[i] = classes.get(label);
|
||||
}
|
||||
return { count, ids, sizes, starts, targets, names, labels };
|
||||
return { count, ids, sizes, starts, targets, names, labels, meta };
|
||||
} finally {
|
||||
closeSync(reader.fd);
|
||||
}
|
||||
}
|
||||
|
||||
function summarize(file) {
|
||||
const { count, ids, sizes, starts, targets, names, labels } = readGraph(file);
|
||||
function summarize(file, keepGraph = false) {
|
||||
const { count, ids, sizes, starts, targets, names, labels, meta } = readGraph(file);
|
||||
const numbers = new Uint32Array(count);
|
||||
const vertices = new Uint32Array(count + 1);
|
||||
const parent = new Uint32Array(count + 1);
|
||||
|
|
@ -300,7 +304,23 @@ function summarize(file) {
|
|||
w = dominator[w];
|
||||
}
|
||||
}
|
||||
return { ids, retained, names, labels, totals, visited, count };
|
||||
if (keepGraph) {
|
||||
parent.fill(count);
|
||||
parent[0] = 0;
|
||||
for (let order = 2; order <= visited; order++) {
|
||||
parent[vertices[order]] = vertices[dominator[order]];
|
||||
}
|
||||
}
|
||||
return {
|
||||
ids,
|
||||
retained,
|
||||
names,
|
||||
labels,
|
||||
totals,
|
||||
visited,
|
||||
count,
|
||||
...(keepGraph ? { starts, targets, immediateDominators: parent, meta } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
function compare(before, after, top) {
|
||||
|
|
@ -380,21 +400,207 @@ function compare(before, after, top) {
|
|||
};
|
||||
}
|
||||
|
||||
try {
|
||||
const args = process.argv.slice(2);
|
||||
const json = args.includes("--json");
|
||||
const files = args.filter((arg) => arg !== "--json");
|
||||
if (files.length !== 2 || files.some((file) => file.startsWith("--"))) {
|
||||
throw new Error(
|
||||
"Usage: node scripts/heap-snapshot-diff.mjs <before.heapsnapshot> <after.heapsnapshot> [--json]",
|
||||
);
|
||||
function readPathEdges(file, meta, wanted) {
|
||||
if (!wanted.size) {
|
||||
return new Map();
|
||||
}
|
||||
const reader = new Reader(file);
|
||||
try {
|
||||
const fields = meta.edge_fields;
|
||||
const typeOffset = fields.indexOf("type");
|
||||
const nameOffset = fields.indexOf("name_or_index");
|
||||
const typeNames = meta.edge_types[typeOffset];
|
||||
const edges = new Map();
|
||||
const neededStrings = new Set();
|
||||
let type;
|
||||
let name;
|
||||
reader.find('"edges"');
|
||||
reader.array((value, index) => {
|
||||
if (!wanted.has(Math.floor(index / fields.length))) {
|
||||
return;
|
||||
}
|
||||
const field = index % fields.length;
|
||||
if (field === typeOffset) {
|
||||
type = typeNames[value];
|
||||
}
|
||||
if (field === nameOffset) {
|
||||
name = value;
|
||||
}
|
||||
if (field === fields.length - 1) {
|
||||
const indexed = type === "element" || type === "hidden";
|
||||
edges.set(Math.floor(index / fields.length), { type, name, indexed });
|
||||
if (!indexed) {
|
||||
neededStrings.add(name);
|
||||
}
|
||||
}
|
||||
});
|
||||
reader.find('"strings"');
|
||||
const strings = new Map();
|
||||
reader.array((value, index) => strings.set(index, value), true, neededStrings);
|
||||
for (const [index, edge] of edges) {
|
||||
edges.set(index, {
|
||||
type: edge.type,
|
||||
name: edge.indexed ? edge.name : strings.get(edge.name),
|
||||
});
|
||||
}
|
||||
return edges;
|
||||
} finally {
|
||||
closeSync(reader.fd);
|
||||
}
|
||||
}
|
||||
|
||||
function retainingPaths(file, snapshot, changes, options) {
|
||||
const { count, starts, targets, ids, names, labels, retained, immediateDominators, meta } =
|
||||
snapshot;
|
||||
const wantedIds = new Set([
|
||||
...changes.dominators.filter((row) => row.delta > 0).map((row) => row.id),
|
||||
...options.nodes,
|
||||
]);
|
||||
const classes = new Map(
|
||||
changes.classes.filter((row) => row.delta > 0).map((row) => [row.label, -1]),
|
||||
);
|
||||
const selected = new Set();
|
||||
for (let node = 0; node < count; node++) {
|
||||
if (wantedIds.delete(ids[node])) {
|
||||
selected.add(node);
|
||||
}
|
||||
const label = labels[names[node]];
|
||||
const previous = classes.get(label);
|
||||
if (previous !== undefined && (previous === -1 || retained[node] > retained[previous])) {
|
||||
classes.set(label, node);
|
||||
}
|
||||
}
|
||||
for (const id of options.nodes) {
|
||||
if (wantedIds.has(id)) {
|
||||
throw new Error(`Node @${id} is absent from the after snapshot`);
|
||||
}
|
||||
}
|
||||
for (const node of classes.values()) {
|
||||
if (node !== -1) {
|
||||
selected.add(node);
|
||||
}
|
||||
}
|
||||
|
||||
// Keep one shortest strong path per node, not all incoming edges or all edge labels.
|
||||
const parents = new Uint32Array(count).fill(count);
|
||||
const parentEdges = new Uint32Array(count);
|
||||
const queue = new Uint32Array(count);
|
||||
parents[0] = 0;
|
||||
const pending = new Set(selected);
|
||||
pending.delete(0);
|
||||
let end = 1;
|
||||
for (let read = 0; read < end && pending.size; read++) {
|
||||
const from = queue[read];
|
||||
for (let edge = starts[from]; edge < starts[from + 1]; edge++) {
|
||||
const to = targets[edge];
|
||||
if (to < count && parents[to] === count) {
|
||||
parents[to] = from;
|
||||
parentEdges[to] = edge;
|
||||
queue[end++] = to;
|
||||
pending.delete(to);
|
||||
}
|
||||
}
|
||||
}
|
||||
const describe = (node) => ({
|
||||
id: ids[node],
|
||||
label: labels[names[node]],
|
||||
retained: retained[node],
|
||||
});
|
||||
const wantedEdges = new Set();
|
||||
function pathTo(node, links, withEdges) {
|
||||
if (links[node] === count) {
|
||||
return null;
|
||||
}
|
||||
const nodes = [];
|
||||
let depth = 0;
|
||||
let current = node;
|
||||
while (true) {
|
||||
if (nodes.length < options.maxDepth) {
|
||||
const entry = describe(current);
|
||||
if (withEdges && current !== 0) {
|
||||
entry.incomingEdge = parentEdges[current];
|
||||
wantedEdges.add(entry.incomingEdge);
|
||||
}
|
||||
nodes.push(entry);
|
||||
}
|
||||
if (current === 0) {
|
||||
break;
|
||||
}
|
||||
current = links[current];
|
||||
depth++;
|
||||
}
|
||||
return { depth, omittedAncestors: depth + 1 - nodes.length, nodes: nodes.toReversed() };
|
||||
}
|
||||
const result = [...selected]
|
||||
.toSorted((a, b) => retained[b] - retained[a])
|
||||
.map((node) =>
|
||||
Object.assign(describe(node), {
|
||||
rootPath: pathTo(node, parents, true),
|
||||
dominatorPath: pathTo(node, immediateDominators, false),
|
||||
}),
|
||||
);
|
||||
// A second streaming pass resolves only displayed edge names, avoiding strings from payloads.
|
||||
const edges = readPathEdges(file, meta, wantedEdges);
|
||||
for (const row of result) {
|
||||
for (const node of row.rootPath?.nodes ?? []) {
|
||||
if (node.incomingEdge !== undefined) {
|
||||
node.incomingEdge = edges.get(node.incomingEdge);
|
||||
}
|
||||
}
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
const usage =
|
||||
"Usage: node scripts/heap-snapshot-diff.mjs <before.heapsnapshot> <after.heapsnapshot> [--json] [--top N] [--node ID] [--max-depth N]";
|
||||
|
||||
function parseArgs(args) {
|
||||
const options = { json: false, files: [], top: 30, nodes: [], maxDepth: 40 };
|
||||
for (let index = 0; index < args.length; index++) {
|
||||
const arg = args[index];
|
||||
if (arg === "--json") {
|
||||
options.json = true;
|
||||
} else if (["--top", "--node", "--max-depth"].includes(arg)) {
|
||||
const value = Number(args[++index]);
|
||||
if (!Number.isSafeInteger(value) || value <= 0) {
|
||||
throw new Error(`${arg} requires a positive safe integer`);
|
||||
}
|
||||
if (arg === "--node") {
|
||||
options.nodes.push(value);
|
||||
} else {
|
||||
options[arg === "--top" ? "top" : "maxDepth"] = value;
|
||||
}
|
||||
} else if (arg.startsWith("--")) {
|
||||
throw new Error(usage);
|
||||
} else {
|
||||
options.files.push(arg);
|
||||
}
|
||||
}
|
||||
if (options.files.length !== 2) {
|
||||
throw new Error(usage);
|
||||
}
|
||||
return options;
|
||||
}
|
||||
|
||||
try {
|
||||
const options = parseArgs(process.argv.slice(2));
|
||||
const { files, json } = options;
|
||||
const before = summarize(files[0]);
|
||||
const after = summarize(files[1]);
|
||||
const after = summarize(files[1], true);
|
||||
const changes = compare(before, after, options.top);
|
||||
const result = {
|
||||
before: { nodes: before.count, reachable: before.visited },
|
||||
after: { nodes: after.count, reachable: after.visited },
|
||||
...compare(before, after, 30),
|
||||
...changes,
|
||||
retainers: retainingPaths(files[1], after, changes, options),
|
||||
notes: [
|
||||
"Node IDs require snapshots from the same process and isolate; this cannot be verified from snapshot contents.",
|
||||
"Retained bytes use the snapshot graph excluding weak and shortcut edges; ephemeron relationships follow V8's encoded internal edges.",
|
||||
"Constructor retained totals overlap between classes; nested instances of one class count once.",
|
||||
"A root path is one shortest strong-edge route, not proof of exclusive ownership. Dominator chains show exclusive retention in this graph, not direct references.",
|
||||
"Retainers include positive top dominator changes, the largest after-node of each positive top constructor, and requested node IDs.",
|
||||
"Truncated paths keep the ancestors nearest the target; omittedAncestors records the missing prefix.",
|
||||
],
|
||||
};
|
||||
if (json) {
|
||||
console.log(JSON.stringify(result, null, 2));
|
||||
|
|
@ -416,6 +622,25 @@ try {
|
|||
);
|
||||
}
|
||||
}
|
||||
console.log("\nRetainer evidence from the after snapshot:");
|
||||
for (const row of result.retainers) {
|
||||
console.log(`\n@${row.id} ${JSON.stringify(row.label)} retained=${row.retained}`);
|
||||
for (const key of ["rootPath", "dominatorPath"]) {
|
||||
const path = row[key];
|
||||
console.log(
|
||||
` ${key}: ${path ? `depth=${path.depth} omittedAncestors=${path.omittedAncestors}` : "unreachable by strong edges"}`,
|
||||
);
|
||||
for (const node of path?.nodes ?? []) {
|
||||
const edge = node.incomingEdge;
|
||||
console.log(
|
||||
` ${edge ? `--${edge.type}:${JSON.stringify(edge.name)}--> ` : ""}@${node.id} ${JSON.stringify(node.label)} retained=${node.retained}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
for (const note of result.notes) {
|
||||
console.log(`\n${note}`);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
console.error(error.message);
|
||||
|
|
|
|||
100
scripts/lib/gateway-heap-rig-catalog.mjs
Normal file
100
scripts/lib/gateway-heap-rig-catalog.mjs
Normal file
|
|
@ -0,0 +1,100 @@
|
|||
import { mkdir, writeFile } from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
|
||||
const PLUGIN_ID = "heap-rig-catalog";
|
||||
|
||||
// The fixture uses the public provider contract and keeps no session cache of its own.
|
||||
const PLUGIN_SOURCE = `import { definePluginEntry } from "openclaw/plugin-sdk/plugin-entry";
|
||||
|
||||
export default definePluginEntry({
|
||||
id: "heap-rig-catalog",
|
||||
name: "Synthetic heap rig catalog",
|
||||
description: "Synthetic catalog projection and publication workload",
|
||||
register(api) {
|
||||
api.registerSessionCatalog({
|
||||
id: "heap-rig-catalog",
|
||||
label: "Synthetic heap rig catalog",
|
||||
audience: "gateway-operators",
|
||||
supportsProcessHomeIsolation: true,
|
||||
async list(params) {
|
||||
if (!params.sessionEntries || !params.agentId) {
|
||||
throw new Error("Synthetic catalog requires the Gateway entry snapshot");
|
||||
}
|
||||
const hostId = "gateway:heap-rig";
|
||||
if (params.hostIds && !params.hostIds.includes(hostId)) {
|
||||
return [];
|
||||
}
|
||||
const search = params.search?.toLowerCase();
|
||||
const rows = params.sessionEntries.entriesForAgent(params.agentId);
|
||||
const sessions = rows
|
||||
.filter(({ sessionKey, entry }) => !search ||
|
||||
(sessionKey + " " + (entry.label ?? "")).toLowerCase().includes(search))
|
||||
.slice(0, params.limitPerHost ?? 100)
|
||||
.map(({ sessionKey, entry }) => ({
|
||||
threadId: entry.sessionId,
|
||||
sessionKey,
|
||||
name: entry.label ?? sessionKey,
|
||||
status: "idle",
|
||||
updatedAt: entry.updatedAt,
|
||||
source: "synthetic-heap-rig",
|
||||
archived: entry.archivedAt !== undefined,
|
||||
canContinue: false,
|
||||
canArchive: false,
|
||||
}));
|
||||
const host = { hostId, label: "Synthetic rig host", kind: "gateway", connected: true, sessions };
|
||||
params.onHost?.(host);
|
||||
if (params.allowPartialResults && params.waitUntil && params.onHost) {
|
||||
// One event-loop turn exercises post-list publication ownership without a timer.
|
||||
const publish = params.onHost;
|
||||
const signal = params.signal;
|
||||
params.waitUntil(new Promise((resolve) => setImmediate(resolve)).then(() => {
|
||||
if (!signal?.aborted) {
|
||||
publish(host);
|
||||
}
|
||||
}));
|
||||
}
|
||||
return [host];
|
||||
},
|
||||
async read({ hostId, threadId }) {
|
||||
return { hostId, threadId, items: [] };
|
||||
},
|
||||
});
|
||||
},
|
||||
});
|
||||
`;
|
||||
|
||||
/** Exercise Gateway catalog lifetimes; this does not emulate a native provider or its caches. */
|
||||
export async function configureHeapRigCatalog(config, root) {
|
||||
const pluginDir = path.join(root, "plugins", PLUGIN_ID);
|
||||
await mkdir(pluginDir, { recursive: true });
|
||||
await Promise.all([
|
||||
writeFile(path.join(pluginDir, "index.mjs"), PLUGIN_SOURCE),
|
||||
writeFile(
|
||||
path.join(pluginDir, "package.json"),
|
||||
JSON.stringify({
|
||||
name: "@openclaw/heap-rig-catalog",
|
||||
private: true,
|
||||
version: "1.0.0",
|
||||
type: "module",
|
||||
openclaw: { extensions: ["./index.mjs"] },
|
||||
}),
|
||||
),
|
||||
writeFile(
|
||||
path.join(pluginDir, "openclaw.plugin.json"),
|
||||
JSON.stringify({
|
||||
id: PLUGIN_ID,
|
||||
activation: { onStartup: true },
|
||||
configSchema: { type: "object", additionalProperties: false },
|
||||
}),
|
||||
),
|
||||
]);
|
||||
const plugins = config.plugins ?? {};
|
||||
config.plugins = {
|
||||
...plugins,
|
||||
enabled: true,
|
||||
...(plugins.allow ? { allow: [...new Set([...plugins.allow, PLUGIN_ID])] } : {}),
|
||||
entries: { ...plugins.entries, [PLUGIN_ID]: { enabled: true } },
|
||||
load: { ...plugins.load, paths: [...(plugins.load?.paths ?? []), pluginDir] },
|
||||
};
|
||||
return { pluginId: PLUGIN_ID, pluginDir };
|
||||
}
|
||||
207
scripts/lib/gateway-heap-rig-turns.mjs
Normal file
207
scripts/lib/gateway-heap-rig-turns.mjs
Normal file
|
|
@ -0,0 +1,207 @@
|
|||
import { randomUUID } from "node:crypto";
|
||||
import { createReadStream } from "node:fs";
|
||||
import { mkdir, rename, stat, writeFile } from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { createInterface } from "node:readline";
|
||||
|
||||
function toolEvents(callId, code) {
|
||||
const item = {
|
||||
type: "function_call",
|
||||
id: `fc_${callId}`,
|
||||
call_id: callId,
|
||||
name: "exec",
|
||||
arguments: JSON.stringify({ title: "Exercise synthetic agent work", code, required: true }),
|
||||
};
|
||||
return [
|
||||
{ type: "response.output_item.added", output_index: 0, item: { ...item, arguments: "" } },
|
||||
{
|
||||
type: "response.function_call_arguments.delta",
|
||||
item_id: item.id,
|
||||
output_index: 0,
|
||||
delta: item.arguments,
|
||||
},
|
||||
{ type: "response.output_item.done", output_index: 0, item },
|
||||
{
|
||||
type: "response.completed",
|
||||
response: {
|
||||
id: `resp_${callId}`,
|
||||
status: "completed",
|
||||
output: [item],
|
||||
usage: { input_tokens: 64, output_tokens: 16, total_tokens: 80 },
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
function assertTerminal(result, runId, marker, tool) {
|
||||
const receipt = result.terminalReceipt;
|
||||
if (
|
||||
result.runId !== runId ||
|
||||
result.status !== "ok" ||
|
||||
result.error ||
|
||||
result.pendingError ||
|
||||
result.yielded ||
|
||||
receipt?.runId !== runId ||
|
||||
!receipt.sessionId ||
|
||||
!receipt.turnId ||
|
||||
receipt.rerouted ||
|
||||
receipt.effective?.provider !== "openai" ||
|
||||
receipt.terminalDisposition !== "visible" ||
|
||||
result.terminalReply?.disposition !== "visible" ||
|
||||
result.terminalReply.text !== marker ||
|
||||
(tool && !receipt.successfulToolNames.includes("exec"))
|
||||
) {
|
||||
throw new Error(`Synthetic turn terminal evidence failed: ${JSON.stringify(result)}`);
|
||||
}
|
||||
return receipt;
|
||||
}
|
||||
|
||||
async function readCodeResult(requestLogPath, start, marker) {
|
||||
const end = (await stat(requestLogPath)).size;
|
||||
if (end <= start) {
|
||||
throw new Error("Mock provider recorded no request for the synthetic tool result");
|
||||
}
|
||||
const input = createReadStream(requestLogPath, { start, end: end - 1 });
|
||||
const lines = createInterface({ input, crlfDelay: Infinity });
|
||||
try {
|
||||
for await (const line of lines) {
|
||||
const record = JSON.parse(line);
|
||||
if (typeof record.body !== "string") {
|
||||
continue;
|
||||
}
|
||||
const body = JSON.parse(record.body);
|
||||
for (const item of body.input ?? []) {
|
||||
if (item.type !== "function_call_output" || typeof item.output !== "string") {
|
||||
continue;
|
||||
}
|
||||
let output;
|
||||
try {
|
||||
output = JSON.parse(item.output);
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
if (output.status === "completed" && output.value?.marker === marker) {
|
||||
return output.value;
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
lines.close();
|
||||
input.destroy();
|
||||
}
|
||||
throw new Error(`Mock provider did not observe completed Code Mode output for ${marker}`);
|
||||
}
|
||||
|
||||
/** Configure the existing mock OpenAI server; the driver owns launch and turn cadence. */
|
||||
export async function configureHeapRigTurns(config, root, mockPort) {
|
||||
await mkdir(root, { recursive: true });
|
||||
const responseControlPath = path.join(root, "mock-turn-responses.json");
|
||||
const requestLogPath = path.join(root, "mock-turn-requests.jsonl");
|
||||
config.tools = { ...config.tools, codeMode: true };
|
||||
const evidence = [];
|
||||
let active = false;
|
||||
|
||||
const writeControl = async (control) => {
|
||||
const temporary = `${responseControlPath}.tmp`;
|
||||
await writeFile(temporary, JSON.stringify(control));
|
||||
await rename(temporary, responseControlPath);
|
||||
};
|
||||
await writeControl({ text: "SYNTHETIC_HEAP_RIG_IDLE" });
|
||||
await writeFile(requestLogPath, "", { flag: "wx" });
|
||||
|
||||
return {
|
||||
responseControlPath,
|
||||
requestLogPath,
|
||||
evidence,
|
||||
env: {
|
||||
MOCK_PORT: String(mockPort),
|
||||
MOCK_BIND_HOST: "127.0.0.1",
|
||||
MOCK_RESPONSE_CONTROL: responseControlPath,
|
||||
MOCK_REQUEST_LOG: requestLogPath,
|
||||
},
|
||||
async runTurn(rpc, index, options = {}) {
|
||||
if (active) {
|
||||
throw new Error("Synthetic heap-rig turns must run serially");
|
||||
}
|
||||
active = true;
|
||||
const startedAt = Date.now();
|
||||
const kind =
|
||||
options.kind ?? (index % 5 === 2 ? "subagent" : index % 5 === 1 ? "code" : "text");
|
||||
const agentId = options.agentId ?? (index % 2 === 0 ? "main" : "research");
|
||||
const sessionKey = options.sessionKey ?? `agent:${agentId}:heap-rig-turn-${index}`;
|
||||
const runId = randomUUID();
|
||||
const marker = `SYNTHETIC_HEAP_RIG_${runId}`;
|
||||
try {
|
||||
const logStart = (await stat(requestLogPath)).size;
|
||||
const code =
|
||||
kind === "subagent"
|
||||
? `const child = await sessions_spawn(${JSON.stringify({
|
||||
task: `Reply with ${marker}.`,
|
||||
runtime: "subagent",
|
||||
context: "isolated",
|
||||
mode: "run",
|
||||
cleanup: "keep",
|
||||
expectsCompletionMessage: false,
|
||||
runTimeoutSeconds: 120,
|
||||
label: `Synthetic heap rig child ${index}`,
|
||||
})});\nreturn { marker: ${JSON.stringify(marker)}, child };`
|
||||
: `const values = Array.from({ length: 256 }, (_, index) => index * index);\nreturn { marker: ${JSON.stringify(marker)}, checksum: values.reduce((sum, value) => sum + value, 0) };`;
|
||||
// The first request consumes the tool fixture. Every continuation and child
|
||||
// gets final text; exact receipt/output checks detect another request stealing it.
|
||||
await writeControl({
|
||||
scriptVersion: runId,
|
||||
responses: kind === "text" ? [{ text: marker }] : [{ events: toolEvents(runId, code) }],
|
||||
default: { text: marker },
|
||||
});
|
||||
if (!options.canStart()) {
|
||||
return null;
|
||||
}
|
||||
const started = await rpc(
|
||||
"agent",
|
||||
{
|
||||
sessionKey,
|
||||
message: `Synthetic heap rig ${kind} turn ${index}. Reply with ${marker}.`,
|
||||
deliver: false,
|
||||
idempotencyKey: runId,
|
||||
},
|
||||
120_000,
|
||||
);
|
||||
if (started.runId !== runId || !["accepted", "ok"].includes(started.status)) {
|
||||
throw new Error(`Synthetic turn was not accepted: ${JSON.stringify(started)}`);
|
||||
}
|
||||
const completed = await rpc("agent.wait", { runId, timeoutMs: 180_000 }, 190_000);
|
||||
const receipt = assertTerminal(completed, runId, marker, kind !== "text");
|
||||
let child;
|
||||
if (kind !== "text") {
|
||||
const value = await readCodeResult(requestLogPath, logStart, marker);
|
||||
if (kind === "subagent") {
|
||||
child = value.child;
|
||||
if (child?.status !== "accepted" || !child.runId || !child.childSessionKey) {
|
||||
throw new Error(`Synthetic child was not accepted: ${JSON.stringify(child)}`);
|
||||
}
|
||||
const childCompleted = await rpc(
|
||||
"agent.wait",
|
||||
{ runId: child.runId, timeoutMs: 180_000 },
|
||||
190_000,
|
||||
);
|
||||
assertTerminal(childCompleted, child.runId, marker, false);
|
||||
}
|
||||
}
|
||||
const result = {
|
||||
index,
|
||||
kind,
|
||||
runId,
|
||||
sessionKey,
|
||||
startedAt,
|
||||
elapsedMs: Date.now() - startedAt,
|
||||
successfulToolNames: receipt.successfulToolNames,
|
||||
...(child ? { childRunId: child.runId, childSessionKey: child.childSessionKey } : {}),
|
||||
};
|
||||
evidence.push(result);
|
||||
return result;
|
||||
} finally {
|
||||
active = false;
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
|
@ -40,12 +40,16 @@ export function resolveGatewayUrl(urlRaw: string): URL {
|
|||
|
||||
export function createGatewayWsClient(params: {
|
||||
url: string;
|
||||
origin?: string;
|
||||
handshakeTimeoutMs?: number;
|
||||
openTimeoutMs?: number;
|
||||
openTimeoutMessage?: string;
|
||||
onEvent?: (evt: GatewayEventFrame) => void;
|
||||
}) {
|
||||
const ws = new WebSocket(params.url, { handshakeTimeout: params.handshakeTimeoutMs ?? 8000 });
|
||||
const ws = new WebSocket(params.url, {
|
||||
handshakeTimeout: params.handshakeTimeoutMs ?? 8000,
|
||||
...(params.origin ? { origin: params.origin } : {}),
|
||||
});
|
||||
ws.binaryType = "nodebuffer";
|
||||
const pending = new Map<
|
||||
string,
|
||||
|
|
|
|||
21
test/scripts/gateway-heap-rig-turns.test.ts
Normal file
21
test/scripts/gateway-heap-rig-turns.test.ts
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
import { afterEach, expect, it, vi } from "vitest";
|
||||
import { configureHeapRigTurns } from "../../scripts/lib/gateway-heap-rig-turns.mjs";
|
||||
import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js";
|
||||
|
||||
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
|
||||
|
||||
it("does not launch a turn when admission closes during asynchronous preparation", async () => {
|
||||
const driver = await configureHeapRigTurns({}, tempDirs.make("heap-rig-turns-"), 19549);
|
||||
const rpc = vi.fn(() => {
|
||||
throw new Error("Expired admission launched a Gateway request");
|
||||
});
|
||||
let admissionOpen = true;
|
||||
const turn = driver.runTurn(rpc, 0, {
|
||||
agentId: "main",
|
||||
canStart: () => admissionOpen,
|
||||
});
|
||||
admissionOpen = false;
|
||||
|
||||
await expect(turn).resolves.toBeNull();
|
||||
expect(rpc).not.toHaveBeenCalled();
|
||||
});
|
||||
|
|
@ -10,7 +10,7 @@ type FixtureNode = {
|
|||
id: number;
|
||||
name: string;
|
||||
size: number;
|
||||
edges: Array<[type: number, target: number]>;
|
||||
edges: Array<[type: number, target: number, name?: string]>;
|
||||
};
|
||||
|
||||
function fixture(grown: boolean) {
|
||||
|
|
@ -26,19 +26,19 @@ function fixture(grown: boolean) {
|
|||
id: 1,
|
||||
name: "root",
|
||||
size: 0,
|
||||
edges: [[0, 1], [0, 2], [1, 7], [2, 4], ...growthEdges],
|
||||
edges: [[0, 1, "cache"], [0, 2], [1, 7], [2, 4, "shortcut"], ...growthEdges],
|
||||
},
|
||||
{
|
||||
id: 3,
|
||||
name: "Cache 🦞",
|
||||
size: 10,
|
||||
edges: [
|
||||
[0, 3],
|
||||
[0, 5],
|
||||
[0, 3, "nested"],
|
||||
[0, 5, "shared"],
|
||||
],
|
||||
},
|
||||
{ id: 5, name: "Cache 🦞", size: 20, edges: [[0, 5]] },
|
||||
{ id: 7, name: "Cache 🦞", size: 5, edges: [[0, 4]] },
|
||||
{ id: 7, name: "Cache 🦞", size: 5, edges: [[0, 4, "payload"]] },
|
||||
{ id: 9, name: "Payload", size: grown ? 90 : 40, edges: [[0, 3]] },
|
||||
{ id: 11, name: "Shared", size: grown ? 70 : 30, edges: [] },
|
||||
{ id: 13, name: "Detached", size: 1_000, edges: [[0, 1]] },
|
||||
|
|
@ -46,7 +46,13 @@ function fixture(grown: boolean) {
|
|||
...growthNodes,
|
||||
];
|
||||
// An unused string crosses the streaming reader's chunk boundary without becoming a class label.
|
||||
const strings = ["x".repeat(1024 * 1024), ...new Set(nodes.map((node) => node.name))];
|
||||
const strings = [
|
||||
"x".repeat(1024 * 1024),
|
||||
...new Set([
|
||||
...nodes.map((node) => node.name),
|
||||
...nodes.flatMap((node) => node.edges.map((edge) => edge[2] ?? "ref")),
|
||||
]),
|
||||
];
|
||||
return JSON.stringify({
|
||||
snapshot: {
|
||||
meta: {
|
||||
|
|
@ -65,21 +71,38 @@ function fixture(grown: boolean) {
|
|||
node.size,
|
||||
node.edges.length,
|
||||
]),
|
||||
edges: nodes.flatMap((node) => node.edges.flatMap(([type, target]) => [type, 0, target * 5])),
|
||||
edges: nodes.flatMap((node) =>
|
||||
node.edges.flatMap(([type, target, name]) => [
|
||||
type,
|
||||
strings.indexOf(name ?? "ref"),
|
||||
target * 5,
|
||||
]),
|
||||
),
|
||||
strings,
|
||||
});
|
||||
}
|
||||
|
||||
it("diffs dominator retention without counting shared, weak, detached, or nested same-class bytes twice", () => {
|
||||
it("diffs exclusive retained sizes and named strong paths without double-counting shared or nested objects", () => {
|
||||
const directory = tempDirs.make("heap-snapshot-diff-");
|
||||
const before = path.join(directory, "before.heapsnapshot");
|
||||
const after = path.join(directory, "after.heapsnapshot");
|
||||
writeFileSync(before, fixture(false));
|
||||
writeFileSync(after, fixture(true));
|
||||
const result = JSON.parse(
|
||||
execFileSync(process.execPath, ["scripts/heap-snapshot-diff.mjs", before, after, "--json"], {
|
||||
encoding: "utf8",
|
||||
}),
|
||||
execFileSync(
|
||||
process.execPath,
|
||||
[
|
||||
"scripts/heap-snapshot-diff.mjs",
|
||||
before,
|
||||
after,
|
||||
"--json",
|
||||
"--max-depth",
|
||||
"3",
|
||||
"--node",
|
||||
"15",
|
||||
],
|
||||
{ encoding: "utf8" },
|
||||
),
|
||||
);
|
||||
|
||||
expect(result.before).toEqual({ nodes: 8, reachable: 6 });
|
||||
|
|
@ -115,4 +138,39 @@ it("diffs dominator retention without counting shared, weak, detached, or nested
|
|||
expect(
|
||||
result.classes.some((row: { label: string }) => /Detached|WeakOnly/u.test(row.label)),
|
||||
).toBe(false);
|
||||
const payload = result.retainers.find((row: { id: number }) => row.id === 9);
|
||||
expect(payload.rootPath).toEqual({
|
||||
depth: 3,
|
||||
omittedAncestors: 1,
|
||||
nodes: [
|
||||
{
|
||||
id: 3,
|
||||
label: "object: Cache 🦞",
|
||||
retained: 105,
|
||||
incomingEdge: { type: "property", name: "cache" },
|
||||
},
|
||||
{
|
||||
id: 7,
|
||||
label: "object: Cache 🦞",
|
||||
retained: 95,
|
||||
incomingEdge: { type: "property", name: "nested" },
|
||||
},
|
||||
{
|
||||
id: 9,
|
||||
label: "object: Payload",
|
||||
retained: 90,
|
||||
incomingEdge: { type: "property", name: "payload" },
|
||||
},
|
||||
],
|
||||
});
|
||||
const shared = result.retainers.find((row: { id: number }) => row.id === 11);
|
||||
expect(shared.rootPath.nodes.map((node: { id: number }) => node.id)).toEqual([1, 3, 11]);
|
||||
expect(shared.dominatorPath.nodes.map((node: { id: number }) => node.id)).toEqual([1, 11]);
|
||||
expect(result.retainers.find((row: { id: number }) => row.id === 15)).toEqual({
|
||||
id: 15,
|
||||
label: "object: WeakOnly",
|
||||
retained: 0,
|
||||
rootPath: null,
|
||||
dominatorPath: null,
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue