mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
perf(plugins): reduce allocation for forwarded stream results (#161696)
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.
This commit is contained in:
parent
d390410b7a
commit
d95272af34
4 changed files with 256 additions and 39 deletions
|
|
@ -128,6 +128,14 @@ its managed handles. Already admitted calls and streams have a bounded chance
|
|||
to finish before disposal; retaining an old function does not make it a current
|
||||
runtime handle.
|
||||
|
||||
Ordinary stream results project their payload on the first `value` read. Nested managed
|
||||
readers share data inspection within that synchronous read, while each reader
|
||||
keeps its own instance admission. An unread terminal payload does not need data
|
||||
inspection. Plain payloads retain their native identity and remain mutable;
|
||||
they are not frozen or transferred. Nested readers recheck later reads for
|
||||
mutations that need executable views. This does not give a closed
|
||||
consumer permission to read a retained active-stream result or call its methods.
|
||||
|
||||
Context engines selected by an admitted turn remain owned through that turn's
|
||||
commit and engine disposal. Replacing an enabled plugin waits for those consumers
|
||||
to close before registering its successor. Disabling or removing a plugin can
|
||||
|
|
|
|||
|
|
@ -6,12 +6,34 @@ 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 }, () =>
|
||||
|
|
@ -21,7 +43,15 @@ const toolResult = {
|
|||
const provider = instance.wrap({
|
||||
async *stream() {
|
||||
for (let sequence = 0; sequence < events; sequence++) {
|
||||
yield { ...payload, sequence };
|
||||
const event =
|
||||
workload === "text"
|
||||
? payload
|
||||
: sequence % 100 === 99
|
||||
? toolResult
|
||||
: sequence % 5 === 4
|
||||
? toolCall
|
||||
: longText;
|
||||
yield { ...event, sequence };
|
||||
}
|
||||
yield toolResult;
|
||||
},
|
||||
|
|
@ -52,13 +82,17 @@ async function consume(layers: number) {
|
|||
}
|
||||
|
||||
try {
|
||||
for (const layers of [0, 5]) {
|
||||
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();
|
||||
|
|
@ -71,21 +105,48 @@ try {
|
|||
await consume(layers);
|
||||
const { profile } = await session.post("HeapProfiler.stopSampling");
|
||||
session.disconnect();
|
||||
const nodes = [profile.head];
|
||||
const nodes = [{ node: profile.head, category: "other" }];
|
||||
const categories: Record<string, number> = {};
|
||||
let bytes = 0;
|
||||
for (const node of nodes) {
|
||||
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;
|
||||
nodes.push(...node.children);
|
||||
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,
|
||||
}),
|
||||
);
|
||||
|
|
|
|||
|
|
@ -21,6 +21,46 @@ afterEach(async () => {
|
|||
});
|
||||
|
||||
describe("plugin async iterable protocol", () => {
|
||||
it.each(["promise", "thenable"] as const)(
|
||||
"awaits a next() %s with native assimilation and live admission",
|
||||
async (kind) => {
|
||||
const instance = owner();
|
||||
const result = { done: false, value: { text: "chunk" } };
|
||||
const receivers: unknown[] = [];
|
||||
const thenable = {
|
||||
// oxlint-disable-next-line unicorn/no-thenable -- Exercise native async-iterator result assimilation.
|
||||
then(resolve: (value: typeof result) => void) {
|
||||
receivers.push(this);
|
||||
expect(instance.hasActiveCall).toBe(true);
|
||||
resolve(result);
|
||||
},
|
||||
};
|
||||
const promise = Promise.resolve(result);
|
||||
// oxlint-disable-next-line unicorn/no-thenable -- Native await must bypass this fulfilled Promise's override.
|
||||
void Object.defineProperty(promise, "then", {
|
||||
get() {
|
||||
throw new Error("Native await must not inspect this Promise's then override");
|
||||
},
|
||||
});
|
||||
const source = instance.wrap({
|
||||
[Symbol.asyncIterator]() {
|
||||
return {
|
||||
next: () => (kind === "promise" ? promise : thenable),
|
||||
return: async () => ({ done: true, value: undefined }),
|
||||
};
|
||||
},
|
||||
});
|
||||
const iterator = source[Symbol.asyncIterator]();
|
||||
try {
|
||||
const next = await iterator.next();
|
||||
expect(next.value).toBe(result.value);
|
||||
expect(receivers).toEqual(kind === "thenable" ? [thenable] : []);
|
||||
} finally {
|
||||
await iterator.return();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["next", "throw"] as const)(
|
||||
"preserves native %s completion after exhaustion while its owner is live",
|
||||
async (method) => {
|
||||
|
|
@ -201,7 +241,7 @@ describe("plugin async iterable protocol", () => {
|
|||
},
|
||||
);
|
||||
|
||||
it("classifies data once per nested read and fences methods added between reads", async () => {
|
||||
it("classifies nested data only when read and fences methods added between reads", async () => {
|
||||
const instance = owner();
|
||||
const inner = instance.retainConsumer();
|
||||
const outer = instance.retainConsumer();
|
||||
|
|
@ -212,12 +252,13 @@ describe("plugin async iterable protocol", () => {
|
|||
})(),
|
||||
);
|
||||
const iterator = outer.wrap(inner.wrap(source))[Symbol.asyncIterator]();
|
||||
const ownKeys = vi.spyOn(Reflect, "ownKeys");
|
||||
try {
|
||||
const next = await iterator.next();
|
||||
if (next.done) {
|
||||
throw new Error("Expected a data chunk");
|
||||
}
|
||||
const ownKeys = vi.spyOn(Reflect, "ownKeys");
|
||||
expect(ownKeys.mock.calls.filter(([value]) => value === payload)).toHaveLength(0);
|
||||
try {
|
||||
expect(next.value).toBe(payload);
|
||||
// Payload traversal must not multiply with the number of managed readers.
|
||||
|
|
@ -232,6 +273,7 @@ describe("plugin async iterable protocol", () => {
|
|||
outer.release();
|
||||
expect(() => changed.read?.()).toThrow(/closed/);
|
||||
} finally {
|
||||
ownKeys.mockRestore();
|
||||
await iterator.return();
|
||||
inner.release();
|
||||
outer.release();
|
||||
|
|
@ -265,6 +307,76 @@ describe("plugin async iterable protocol", () => {
|
|||
}
|
||||
});
|
||||
|
||||
it("keeps native result exemptions from admitting a nested callable payload", async () => {
|
||||
const instance = owner();
|
||||
const outer = instance.retainConsumer();
|
||||
const payload = { read: () => "chunk" };
|
||||
const result = Object.assign(Object.setPrototypeOf(new Date(), Object.prototype), {
|
||||
done: false,
|
||||
value: payload,
|
||||
});
|
||||
const source = instance.wrap({
|
||||
[Symbol.asyncIterator]() {
|
||||
return {
|
||||
next: async () => result,
|
||||
return: async () => ({ done: true, value: undefined }),
|
||||
};
|
||||
},
|
||||
});
|
||||
const iterator = outer.wrap(source)[Symbol.asyncIterator]();
|
||||
try {
|
||||
const next = await iterator.next();
|
||||
const chunk = next.value;
|
||||
expect(chunk?.read()).toBe("chunk");
|
||||
await iterator.return();
|
||||
outer.release();
|
||||
expect(() => chunk?.read()).toThrow(/closed/);
|
||||
} finally {
|
||||
await iterator.return();
|
||||
outer.release();
|
||||
}
|
||||
});
|
||||
|
||||
it.each(["record", "callable", "proxy"] as const)(
|
||||
"preserves unread terminal %s results after self-retirement",
|
||||
async (kind) => {
|
||||
const instance = owner();
|
||||
const payload = { content: [{ text: "complete" }] };
|
||||
const result = { done: true, value: payload };
|
||||
const callable = Object.assign(() => {}, result);
|
||||
Object.defineProperty(callable, "length", {
|
||||
get() {
|
||||
expect(instance.hasActiveCall).toBe(true);
|
||||
return 0;
|
||||
},
|
||||
});
|
||||
const source = instance.wrap({
|
||||
[Symbol.asyncIterator]() {
|
||||
return {
|
||||
async next() {
|
||||
await instance.dispose();
|
||||
return kind === "callable"
|
||||
? callable
|
||||
: kind === "proxy"
|
||||
? Object.create(new Proxy(result, {}))
|
||||
: result;
|
||||
},
|
||||
};
|
||||
},
|
||||
});
|
||||
const iterator = source[Symbol.asyncIterator]();
|
||||
const next = await iterator.next();
|
||||
expect(next.done).toBe(true);
|
||||
await instance.dispose();
|
||||
if (kind === "proxy") {
|
||||
expect(() => next.value).toThrow(/reloaded or disabled/);
|
||||
} else {
|
||||
expect(next.value).toBe(payload);
|
||||
expect(structuredClone(next.value)).toEqual(payload);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it("admits a replacement getter on an otherwise core-owned iterator result", async () => {
|
||||
const instance = owner();
|
||||
const original = instance.wrap(
|
||||
|
|
|
|||
|
|
@ -108,6 +108,16 @@ function hasProxyPrototype(object: object): boolean {
|
|||
return false;
|
||||
}
|
||||
|
||||
function isNativePluginData(value: object): boolean {
|
||||
return (
|
||||
types.isAnyArrayBuffer(value) ||
|
||||
types.isArrayBufferView(value) ||
|
||||
types.isDate(value) ||
|
||||
types.isRegExp(value) ||
|
||||
types.isNativeError(value)
|
||||
);
|
||||
}
|
||||
|
||||
function isPluginData(
|
||||
value: unknown,
|
||||
seen?: Set<object>,
|
||||
|
|
@ -120,13 +130,7 @@ function isPluginData(
|
|||
if (types.isProxy(value)) {
|
||||
return false;
|
||||
}
|
||||
if (
|
||||
types.isAnyArrayBuffer(value) ||
|
||||
types.isArrayBufferView(value) ||
|
||||
types.isDate(value) ||
|
||||
types.isRegExp(value) ||
|
||||
types.isNativeError(value)
|
||||
) {
|
||||
if (isNativePluginData(value)) {
|
||||
return true;
|
||||
}
|
||||
if (seen?.has(value)) {
|
||||
|
|
@ -293,14 +297,10 @@ export function createPluginValueView(
|
|||
wrap: (value) => wrap(value),
|
||||
invoke: admitCallback,
|
||||
});
|
||||
const wrapResult = <T>(
|
||||
result: T,
|
||||
callerData?: unknown[],
|
||||
project: <V>(value: V) => V = wrap,
|
||||
): T => {
|
||||
const wrapResult = <T>(result: T, callerData?: unknown[]): T => {
|
||||
const completion = resolvePluginReturnPromise(result);
|
||||
if (completion) {
|
||||
const pending = mapPluginReturnPromise(completion, (resolved) => project(resolved));
|
||||
const pending = mapPluginReturnPromise(completion, (resolved) => wrap(resolved));
|
||||
if (pending.host) {
|
||||
valueInstances.setHost(pending.value, bindings.instance);
|
||||
} else {
|
||||
|
|
@ -309,7 +309,7 @@ export function createPluginValueView(
|
|||
// SAFETY: Promise-like results retain their resolved type while callable values stay owned.
|
||||
return pending.value as T;
|
||||
}
|
||||
return callerData?.includes(result) ? result : project(result);
|
||||
return callerData?.includes(result) ? result : wrap(result);
|
||||
};
|
||||
|
||||
/** Callables retain their instance; schemas remain data for host validators. */
|
||||
|
|
@ -571,14 +571,28 @@ export function createPluginValueView(
|
|||
};
|
||||
|
||||
// Host readers retain their lease; the outer iterator can project their payload lazily.
|
||||
const wrapIteratorResult = <T>(value: T): T =>
|
||||
value !== null &&
|
||||
typeof value === "object" &&
|
||||
IteratorResultReader.get(value) &&
|
||||
!pluginMemberNeedsAdmission(value, "then") &&
|
||||
typeof Reflect.get(value, "then") !== "function"
|
||||
? value
|
||||
: wrap(value);
|
||||
const wrapIteratorResult = (value: object, dataRead?: IteratorDataRead): object => {
|
||||
const result =
|
||||
IteratorResultReader.get(value) &&
|
||||
!pluginMemberNeedsAdmission(value, "then") &&
|
||||
typeof Reflect.get(value, "then") !== "function"
|
||||
? value
|
||||
: wrap(value);
|
||||
if (
|
||||
dataRead &&
|
||||
result === value &&
|
||||
!types.isProxy(value) &&
|
||||
!IteratorResultReader.get(value) &&
|
||||
!isNativePluginData(value)
|
||||
) {
|
||||
// Only a complete graph inspection admits the payload for this synchronous read.
|
||||
const payload: unknown = Object.getOwnPropertyDescriptor(value, "value")?.value;
|
||||
if (payload !== null && typeof payload === "object") {
|
||||
dataRead.data = payload;
|
||||
}
|
||||
}
|
||||
return result;
|
||||
};
|
||||
|
||||
const admitIterator = (iterator: object): PluginIteratorAdmission => {
|
||||
const current = iterators.get(iterator);
|
||||
|
|
@ -690,22 +704,44 @@ export function createPluginValueView(
|
|||
}
|
||||
throw new TypeError("Plugin iterator method must be callable");
|
||||
}
|
||||
const next: unknown = await wrapResult(
|
||||
Reflect.apply(method, iterator, args),
|
||||
undefined,
|
||||
wrapIteratorResult,
|
||||
);
|
||||
const next: unknown = await Reflect.apply(method, iterator, args);
|
||||
if (next === null || (typeof next !== "object" && typeof next !== "function")) {
|
||||
throw new TypeError("Plugin async iterator result must be an object");
|
||||
}
|
||||
const complete = Boolean(readResultMember(next, "done"));
|
||||
const doneDescriptor =
|
||||
typeof next === "object" &&
|
||||
!hasProxyPrototype(next) &&
|
||||
Object.getOwnPropertyDescriptor(next, "done");
|
||||
// Executable reflection and non-Boolean completion values keep their original projection.
|
||||
let projected =
|
||||
doneDescriptor &&
|
||||
"value" in doneDescriptor &&
|
||||
typeof doneDescriptor.value === "boolean" &&
|
||||
!IteratorResultReader.get(next)
|
||||
? undefined
|
||||
: wrapIteratorResult(next);
|
||||
const complete = Boolean(readResultMember(projected ?? next, "done"));
|
||||
// IteratorClose ends this admission even when a generator yields in finally.
|
||||
// A later explicit next can acquire a new lease only while the instance is live.
|
||||
state = complete ? "done" : key === "return" ? "returned" : state;
|
||||
const readWithData = (dataRead: IteratorDataRead) =>
|
||||
active
|
||||
? readResultMember(next, "value", dataRead)
|
||||
: Reflect.get(IteratorResultReader.get(next) ? wrap(next) : next, "value");
|
||||
const readWithData = (dataRead: IteratorDataRead) => {
|
||||
if (!projected) {
|
||||
const prepare = () => {
|
||||
projected = wrapIteratorResult(next, dataRead);
|
||||
};
|
||||
// Result projection belongs to its first payload read, not next()'s completion.
|
||||
// Terminal data stays readable after the iterator's admission has finished.
|
||||
if (active) {
|
||||
invoke(prepare);
|
||||
} else {
|
||||
prepare();
|
||||
}
|
||||
}
|
||||
const result = projected!;
|
||||
return active
|
||||
? readResultMember(result, "value", dataRead)
|
||||
: Reflect.get(IteratorResultReader.get(result) ? wrap(result) : result, "value");
|
||||
};
|
||||
const readValue = () => readWithData({});
|
||||
// Completion may join disposal; value stays lazy and checks the exact inner lease.
|
||||
const result = Object.defineProperty({ done: complete }, "value", {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue