mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 17:53:39 +00:00
Project ordinary iterator payloads on first consumption and share completed graph inspection only within that synchronous reader chain. Preserve consumer admission, mutable native data, and exceptional-result projection while removing the intermediate Promise mapping. Reduce sampled allocation by 38-40% across five forwarding layers. Record the read-before-forward limitation, add native-await and lifecycle regressions, and clarify the SDK contract.
156 lines
5.2 KiB
TypeScript
156 lines
5.2 KiB
TypeScript
// Run with node --expose-gc --import ./scripts/tsx.mjs scripts/bench-plugin-stream.ts.
|
|
import { Session } from "node:inspector/promises";
|
|
import { performance } from "node:perf_hooks";
|
|
import { PluginInstance } from "../src/plugins/plugin-instance.js";
|
|
import { createEmptyPluginRegistry } from "../src/plugins/registry-empty.js";
|
|
import { createPluginRecord } from "../src/plugins/status.test-helpers.js";
|
|
|
|
const events = 2_000;
|
|
const workload = process.argv[2] ?? "text";
|
|
if (workload !== "text" && workload !== "mixed") {
|
|
throw new Error("Expected text or mixed workload");
|
|
}
|
|
const registry = createEmptyPluginRegistry();
|
|
const record = createPluginRecord({ id: "stream-fixture", origin: "bundled" });
|
|
registry.plugins.push(record);
|
|
const instance = new PluginInstance(record.id, { record, registry });
|
|
const content = Array.from({ length: 64 }, () => ({ type: "text", text: "x".repeat(38) }));
|
|
const payload = { type: "text_delta", delta: "chunk", partial: { content } };
|
|
const longText = {
|
|
...payload,
|
|
partial: { content: [...content, { type: "text", text: "assistant text ".repeat(2_000) }] },
|
|
};
|
|
const toolCall = {
|
|
type: "toolcall_delta",
|
|
delta: '"argument":',
|
|
partial: {
|
|
content: [
|
|
{
|
|
type: "toolCall",
|
|
id: "call_fixture",
|
|
name: "read",
|
|
arguments: { text: "x".repeat(8_000) },
|
|
},
|
|
],
|
|
},
|
|
};
|
|
const toolResult = {
|
|
type: "tool_result",
|
|
content: Array.from({ length: 50 }, () =>
|
|
content.map((block) => ({ type: block.type, text: block.text })),
|
|
),
|
|
};
|
|
const provider = instance.wrap({
|
|
async *stream() {
|
|
for (let sequence = 0; sequence < events; sequence++) {
|
|
const event =
|
|
workload === "text"
|
|
? payload
|
|
: sequence % 100 === 99
|
|
? toolResult
|
|
: sequence % 5 === 4
|
|
? toolCall
|
|
: longText;
|
|
yield { ...event, sequence };
|
|
}
|
|
yield toolResult;
|
|
},
|
|
});
|
|
function transform(source: AsyncIterable<unknown>): AsyncIterable<unknown> {
|
|
return instance.wrap({
|
|
[Symbol.asyncIterator]() {
|
|
const iterator = source[Symbol.asyncIterator]();
|
|
return {
|
|
next: () => iterator.next(),
|
|
return: async () => iterator.return?.() ?? { done: true, value: undefined },
|
|
};
|
|
},
|
|
});
|
|
}
|
|
async function consume(layers: number) {
|
|
let source: AsyncIterable<unknown> = provider.stream();
|
|
for (let layer = 0; layer < layers; layer++) {
|
|
source = transform(source);
|
|
}
|
|
let count = 0;
|
|
for await (const _ of source) {
|
|
count++;
|
|
}
|
|
if (count !== events + 1) {
|
|
throw new Error(`Lost events: ${count}`);
|
|
}
|
|
}
|
|
|
|
try {
|
|
for (const layers of [0, 1, 5]) {
|
|
await consume(layers);
|
|
const samples: number[] = [];
|
|
const cpuSamples: number[] = [];
|
|
for (let sample = 0; sample < 5; sample++) {
|
|
const start = performance.now();
|
|
const cpu = process.threadCpuUsage();
|
|
await consume(layers);
|
|
samples.push((performance.now() - start) / (events + 1));
|
|
const used = process.threadCpuUsage(cpu);
|
|
cpuSamples.push((used.user + used.system) / 1_000 / (events + 1));
|
|
}
|
|
globalThis.gc?.();
|
|
const session = new Session();
|
|
session.connect();
|
|
await session.post("HeapProfiler.startSampling", {
|
|
samplingInterval: 512,
|
|
includeObjectsCollectedByMajorGC: true,
|
|
includeObjectsCollectedByMinorGC: true,
|
|
});
|
|
await consume(layers);
|
|
const { profile } = await session.post("HeapProfiler.stopSampling");
|
|
session.disconnect();
|
|
const nodes = [{ node: profile.head, category: "other" }];
|
|
const categories: Record<string, number> = {};
|
|
let bytes = 0;
|
|
for (const entry of nodes) {
|
|
const { node } = entry;
|
|
const name = node.callFrame.functionName;
|
|
const category =
|
|
name === "isPluginData"
|
|
? "dataWalk"
|
|
: entry.category === "dataWalk"
|
|
? entry.category
|
|
: name === "structuredClone"
|
|
? "structuredClone"
|
|
: name === "wrap"
|
|
? "valueViews"
|
|
: name === "wrapIteratorResult" || name === "IteratorResultReader"
|
|
? "iteratorResults"
|
|
: name === "readResultMember"
|
|
? "resultReads"
|
|
: entry.category;
|
|
bytes += node.selfSize;
|
|
categories[category] = (categories[category] ?? 0) + node.selfSize;
|
|
nodes.push(...node.children.map((child) => ({ node: child, category })));
|
|
}
|
|
console.log(
|
|
JSON.stringify({
|
|
node: process.version,
|
|
events,
|
|
workload,
|
|
layers,
|
|
eventBytes: Buffer.byteLength(JSON.stringify(payload)),
|
|
toolResultBytes: Buffer.byteLength(JSON.stringify(toolResult)),
|
|
msPerEvent: samples.toSorted((a, b) => a - b)[2],
|
|
mainThreadMsPer10kEvents: cpuSamples.toSorted((a, b) => a - b)[2]! * 10_000,
|
|
sampledAllocatedBytesPerEvent: bytes / (events + 1),
|
|
sampledMBPer10kEvents: bytes / (events + 1) / 100,
|
|
allocationCategories: Object.fromEntries(
|
|
Object.entries(categories).map(([name, size]) => [
|
|
name,
|
|
{ bytesPerEvent: size / (events + 1), percent: (100 * size) / bytes },
|
|
]),
|
|
),
|
|
rssBytes: process.memoryUsage().rss,
|
|
}),
|
|
);
|
|
}
|
|
} finally {
|
|
await instance.dispose();
|
|
}
|