diff --git a/config/knip.config.ts b/config/knip.config.ts index 41c650d11871..a994e17c2c58 100644 --- a/config/knip.config.ts +++ b/config/knip.config.ts @@ -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. diff --git a/docs/gateway/diagnostics.md b/docs/gateway/diagnostics.md index 5a1d3c57f824..162243ccbe1a 100644 --- a/docs/gateway/diagnostics.md +++ b/docs/gateway/diagnostics.md @@ -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 ` 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 diff --git a/scripts/gateway-heap-rig.mjs b/scripts/gateway-heap-rig.mjs new file mode 100644 index 000000000000..126ff3ecf8f6 --- /dev/null +++ b/scripts/gateway-heap-rig.mjs @@ -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} */ + 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} */ + 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); +} diff --git a/scripts/heap-snapshot-diff.mjs b/scripts/heap-snapshot-diff.mjs index 303bc93c4af3..3c21959a66cf 100644 --- a/scripts/heap-snapshot-diff.mjs +++ b/scripts/heap-snapshot-diff.mjs @@ -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 [--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 [--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); diff --git a/scripts/lib/gateway-heap-rig-catalog.mjs b/scripts/lib/gateway-heap-rig-catalog.mjs new file mode 100644 index 000000000000..04d53b6b5c5e --- /dev/null +++ b/scripts/lib/gateway-heap-rig-catalog.mjs @@ -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 }; +} diff --git a/scripts/lib/gateway-heap-rig-turns.mjs b/scripts/lib/gateway-heap-rig-turns.mjs new file mode 100644 index 000000000000..ee90b5d08fe6 --- /dev/null +++ b/scripts/lib/gateway-heap-rig-turns.mjs @@ -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; + } + }, + }; +} diff --git a/scripts/lib/gateway-ws-client.ts b/scripts/lib/gateway-ws-client.ts index 16cedbdb6b98..8a313472cfb1 100644 --- a/scripts/lib/gateway-ws-client.ts +++ b/scripts/lib/gateway-ws-client.ts @@ -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, diff --git a/test/scripts/gateway-heap-rig-turns.test.ts b/test/scripts/gateway-heap-rig-turns.test.ts new file mode 100644 index 000000000000..f99c5d6d92b7 --- /dev/null +++ b/test/scripts/gateway-heap-rig-turns.test.ts @@ -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(); +}); diff --git a/test/scripts/heap-snapshot-diff.test.ts b/test/scripts/heap-snapshot-diff.test.ts index c73373c9ab79..e23a7dabd75f 100644 --- a/test/scripts/heap-snapshot-diff.test.ts +++ b/test/scripts/heap-snapshot-diff.test.ts @@ -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, + }); });