openclaw/scripts/bench-usage-refresh-memory.ts
Peter Steinberger 5b90a73498
fix(usage): prevent refresh OOMs on large SQLite sessions (#156937)
Closes #156898

## What Problem This Solves

Large SQLite sessions can repeatedly exhaust the usage-refresh worker's 512 MiB heap, leaving usage totals missing or stale after an otherwise successful `sessions.usage` response.

## User Impact

Large sessions can finish refreshing without increasing the worker limit. Doctor reports bounded, per-session refresh failures and successful refreshes clear their warnings. No configuration, rollup format, database schema, or migration changes are required. Thanks @Conan-Scott for the detailed report and allocation control.

## Why This Change Was Made

The reader eagerly decoded the entire requested range before aggregation's 128-record batches. It now reads at most 1,024 rows and 8 MiB of decoded JSON per page, with one lookahead event (an individual oversized identity event is still accepted). A first pass retains compact navigation facts; selected bodies then feed the existing aggregation in ancestry order. Append validation, reset/leaf fallback, checkpoint checks, and conditional publication remain in place. Navigation metadata still scales with event count; payload bodies do not. The existing incognito host-frame reader is unchanged.

Refresh failures use the existing asynchronous core plugin-state storage, capped at 256 session facts, without storing raw error or transcript content. The cache owner clears each fact after successful publication; Doctor displays unresolved failures.

## Evidence

All heavy validation ran on Blacksmith Testbox (`blacksmith-testbox`, profile `openclaw-check`). Primary measurements and typechecks: lease `tbx_01m38g2xckkz1apk1smezhftg1`, Node 24.19.0, [run](https://github.com/openclaw/openclaw/actions/runs/35942606687). Final lint, guards, paging tests, and benchmark replay: lease `tbx_01m38k1hj8p9r5ay910z20s2sv`, [run](https://github.com/openclaw/openclaw/actions/runs/35946363096). Both leases stopped after validation.

The opt-in `node --import tsx scripts/bench-usage-refresh-memory.ts` fixture uses 24,000 events, 6,000 identity rows and 18,000 production zstd rows: 1,186,945,771 decoded bytes (1.105 GiB), with each event below 4 MiB. Both runs use the real refresh worker and its production 512 MiB limit; the baseline substitutes the original reader/scanner from `57f912b5f7` before executing the same harness with `--expect-oom`.

| Measurement | Original reader | Bounded reader |
| --- | --- | --- |
| Outcome | `ERR_WORKER_OUT_OF_MEMORY` | Complete, exact reference rollup |
| Sampled peak worker heap | 489.8 MiB | 228.9 MiB (53.3% lower) |
| Refresh duration | Failed after 2.39 s | Completed in 4.64 s |

The final-source replay completed in 4.40 s at 234.4 MiB sampled peak, again with exact rollup equality and below the asserted 384 MiB bound.

Heap sampling uses `Worker.getHeapStatistics()` every 10 ms; these are observed peaks, not complete allocation profiles. Failure time is not a throughput comparison. The complete rollup matched the bounded JSONL worker reference, independently checked for 24,000 records, 240,000 tokens, and 24,000 synthetic cost units. Rollup SHA-256: `f0f184c9b958a74211b6d66678bb1808497d07a96b6b8a84cfb1caabfdc9a856`.

- `node scripts/run-vitest.mjs src/infra/session-cost-usage-worker-refresh.test.ts src/infra/session-cost-usage-cache.worker.test.ts src/commands/doctor-usage-cost-cache.test.ts src/infra/session-cost-usage-worker-io.test.ts --maxWorkers=1`: 27 tests passed, 65.98 s total including cold worker compilation. New paging file: 5 tests / 7 ms. Changed worker file: 8 tests / 14.93 s; new failure→Doctor→recovery test: 1.76 s.
- `node scripts/check-changed.mjs`: all selected checks completed successfully across the initial run and resumed native-plan commands. Core types and all 25 test-typecheck shards passed; script/root-test types, lint, formatting, storage/import-cycle/security guards also passed. The initial run caught optional `Worker.resourceLimits` typing; lint caught parameter reassignment and benchmark brace/untyped-throw issues. All were corrected and their failed checks replayed successfully.
- Final `node scripts/run-vitest.mjs src/infra/session-cost-usage-worker-refresh.test.ts --maxWorkers=1`: 5 passed / 7 ms test time, 22.28 s wrapper wall including cold worker compilation.
- Final `node --import tsx scripts/bench-usage-refresh-memory.ts`: passed, same rollup hash as the paired measurement.
- Independent Codex review: no actionable P0–P2 findings.

Initial proof attempts exposed two fixture/setup issues, both corrected before collecting the measurements: missing Testbox checkout dependencies and a synthetic session row that bypassed canonical admission. The final fixture creates the session through its owner. Two later native checksum sync attempts timed out before validation; publishing the reviewed branch and warming a fresh lease recovered native targeted sync. No production data or Gateway deployment was used; this proves the synthetic allocation defect and rollup equality, not attribution of all reported Gateway RSS growth.
2026-09-24 02:50:28 +00:00

314 lines
11 KiB
TypeScript

// Opt-in, synthetic real-worker regression. Run on a Testbox with Node 24+:
// pnpm bench:usage-refresh-memory [--expect-oom]
import assert from "node:assert/strict";
import { createHash } from "node:crypto";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { parseArgs } from "node:util";
import type { Worker } from "node:worker_threads";
import { parseStrictIntegerOption } from "./lib/strict-integer-option.ts";
const { values } = parseArgs({
options: {
events: { type: "string", default: "24000" },
"payload-kib": { type: "string", default: "48" },
"max-heap-mib": { type: "string", default: "384" },
"expect-oom": { type: "boolean", default: false },
},
});
const count = parseStrictIntegerOption({
raw: values.events,
fallback: 24000,
min: 1,
label: "--events",
});
const payloadBytes =
1024 *
parseStrictIntegerOption({
raw: values["payload-kib"],
fallback: 48,
min: 1,
label: "--payload-kib",
});
const maxHeapBytes =
1024 *
1024 *
parseStrictIntegerOption({
raw: values["max-heap-mib"],
fallback: 384,
min: 1,
label: "--max-heap-mib",
});
assert.ok(payloadBytes < 4 * 1024 * 1024 - 1024, "Each event must fit the production 4 MiB limit");
const root = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-usage-memory-"));
process.env.OPENCLAW_STATE_DIR = root;
process.env.OPENCLAW_CONFIG_PATH = path.join(root, "openclaw.json");
fs.writeFileSync(process.env.OPENCLAW_CONFIG_PATH, "{}\n");
const agentId = "usage-memory-benchmark";
const sessionId = "synthetic-large-session";
const timestamp = "2026-09-23T12:00:00.000Z";
const referenceFile = path.join(root, "reference.jsonl");
type HeapSample = {
threadId: number;
limitMiB: number | undefined;
peakHeapBytes: number;
samples: number;
pending: boolean;
};
const heaps = new Map<Worker, HeapSample>();
let measuring = false;
function observeWorker(worker: Worker): void {
if (!measuring) {
return;
}
heaps.set(worker, {
threadId: worker.threadId,
limitMiB: worker.resourceLimits?.maxOldGenerationSizeMb,
peakHeapBytes: 0,
samples: 0,
pending: false,
});
}
process.on("worker", observeWorker);
function sampleHeaps(): void {
for (const [worker, sample] of heaps) {
if (sample.pending || worker.threadId < 0) {
continue;
}
sample.pending = true;
void worker
.getHeapStatistics()
.then((heap) => {
sample.peakHeapBytes = Math.max(sample.peakHeapBytes, heap.used_heap_size);
sample.samples++;
})
.catch(() => {})
.finally(() => {
sample.pending = false;
});
}
}
function errorChain(error: unknown): string {
if (!(error instanceof Error)) {
return String(error);
}
const code = Reflect.get(error, "code");
return `${error.name}: ${code ?? ""} ${error.message}${error.cause ? `; ${errorChain(error.cause)}` : ""}${error instanceof AggregateError ? error.errors.map(errorChain).join("; ") : ""}`;
}
const { openOpenClawAgentDatabase, closeOpenClawAgentDatabasesAsync } =
await import("../src/state/openclaw-agent-db.js");
const { closeOpenClawStateDatabaseAsync } = await import("../src/state/openclaw-state-db.js");
const { prepareTranscriptPayload } = await import("../src/config/sessions/transcript-payload.js");
const { createSessionEntryWithTranscript } =
await import("../src/config/sessions/session-accessor.entry-mutation.js");
const { waitForSessionTranscriptIndexReconcilesInStateDir } =
await import("../src/config/sessions/session-transcript-reconcile.js");
const { closeSessionTranscriptReconcileWorkerPool } =
await import("../src/config/sessions/session-transcript-reconcile-pool.js");
const { formatSqliteSessionFileMarker } =
await import("../src/config/sessions/legacy-sqlite-marker.js");
const { prepareUsageCostWorker, runUsageCostWorker } =
await import("../src/infra/session-cost-usage-worker-runtime.js");
const { decodeUsageCostRollup, decodeUsageCostRollupEnvelope, USAGE_COST_ROLLUP_SCOPE } =
await import("../src/infra/session-cost-usage-rollup-codec.js");
const { rotateDatabaseWorkers, costRefreshLane } =
await import("../src/config/sessions/session-transcript-worker-resources.js");
try {
const database = openOpenClawAgentDatabase({ agentId });
const sessionKey = `agent:${agentId}:benchmark`;
const created = await createSessionEntryWithTranscript(
{ agentId, sessionKey, storePath: database.path },
() => ({ ok: true, entry: { sessionId, updatedAt: Date.parse(timestamp) } }),
{ cwd: root },
);
assert.ok(created.ok);
await waitForSessionTranscriptIndexReconcilesInStateDir(root);
database.db.prepare("DELETE FROM transcript_events WHERE session_id = ?").run(sessionId);
const insert = database.db.prepare(
"INSERT INTO transcript_events (session_id, seq, event_json, event_zstd, event_utf8_bytes, navigation_json, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)",
);
let decodedBytes = 0;
let identityRows = 0;
let compressedRows = 0;
const referenceFd = fs.openSync(referenceFile, "w");
const fixtureStarted = performance.now();
try {
// Fixture-only direct writes keep setup bounded, with no live operator state.
database.db.exec("BEGIN");
for (let start = 0; start < count; start += 128) {
const referenceBatch: string[] = [];
for (let index = start; index < Math.min(count, start + 128); index++) {
const event = {
type: "message",
id: `event-${index}`,
parentId: index === 0 ? null : `event-${index - 1}`,
timestamp,
message: {
role: "assistant",
provider: "synthetic",
model: "fixture",
content: [{ type: "text", text: "Synthetic benchmark message" }],
usage: { input: 7, output: 3, totalTokens: 10, cost: { total: 1 } },
},
};
// Text is irrelevant to the usage rollup; the independent JSONL backend
// supplies the same contribution without retaining the padded payload.
referenceBatch.push(JSON.stringify(event));
event.message.content[0]!.text += "x".repeat(payloadBytes);
const json = JSON.stringify(event);
decodedBytes += Buffer.byteLength(json);
const payload =
index % 4 === 0
? {
event_json: json,
event_zstd: null,
event_utf8_bytes: Buffer.byteLength(json),
navigation_json: null,
}
: prepareTranscriptPayload(database.db, json, event);
if (payload.event_zstd) {
compressedRows++;
} else {
identityRows++;
}
insert.run(
sessionId,
index + 1,
payload.event_json,
payload.event_zstd,
payload.event_utf8_bytes,
payload.navigation_json,
Date.parse(timestamp),
);
}
fs.writeSync(referenceFd, `${referenceBatch.join("\n")}\n`);
}
database.db.exec("COMMIT");
} catch (error) {
database.db.exec("ROLLBACK");
throw error;
} finally {
fs.closeSync(referenceFd);
}
assert.ok(
identityRows > 0 && compressedRows > 0,
"Fixture must exercise identity and compressed payloads",
);
console.log(
JSON.stringify({
phase: "fixture",
events: count,
decodedBytes,
identityRows,
compressedRows,
fixtureMs: performance.now() - fixtureStarted,
}),
);
const marker = formatSqliteSessionFileMarker({ agentId, sessionId, storePath: database.path });
const prepared = prepareUsageCostWorker({
agentId,
databasePath: database.path,
storePath: database.path,
sessionFiles: [marker, referenceFile],
});
function readRollup(key: string) {
const row = database.db
.prepare("SELECT value_json, blob FROM cache_entries WHERE scope = ? AND key = ?")
.get(USAGE_COST_ROLLUP_SCOPE, key);
assert.ok(
row && typeof row.value_json === "string" && row.blob instanceof Uint8Array,
"Worker must persist a rollup",
);
const envelope = decodeUsageCostRollupEnvelope(row.value_json);
assert.ok(envelope);
const entry = decodeUsageCostRollup(row.value_json, envelope.pricingFingerprint, row.blob);
assert.ok(entry);
return entry;
}
assert.equal(
(await runUsageCostWorker(prepared, { kind: "refresh", sessionFiles: [referenceFile] })).kind,
"refresh",
);
const reference = readRollup(referenceFile);
assert.equal(reference.parsedRecords, count);
assert.equal(reference.countedRecords, count);
const bucket = Object.values(reference.rollup.buckets)[0];
assert.equal(bucket?.totals.totalTokens, count * 10);
assert.equal(bucket?.totals.totalCost, count);
// Start a fresh production refresh worker so its whole lifetime is observed.
await rotateDatabaseWorkers(costRefreshLane);
measuring = true;
const sampler = setInterval(sampleHeaps, 10);
const started = performance.now();
let failure: unknown;
try {
assert.equal(
(await runUsageCostWorker(prepared, { kind: "refresh", sessionFiles: [marker] })).kind,
"refresh",
);
} catch (error) {
failure = error;
} finally {
measuring = false;
clearInterval(sampler);
}
const refreshMs = performance.now() - started;
const workers = [...heaps.values()].map(({ pending: _pending, ...sample }) => sample);
const peakHeapBytes = Math.max(
0,
...workers.filter((worker) => worker.limitMiB === 512).map((worker) => worker.peakHeapBytes),
);
console.log(
JSON.stringify({
phase: "refresh",
refreshMs,
peakHeapBytes,
heapBoundBytes: maxHeapBytes,
workers,
outcome: failure ? errorChain(failure) : "success",
sampling: "Node Worker.getHeapStatistics every 10 ms; observed peak, not an allocation total",
}),
);
assert.ok(peakHeapBytes > 0, "Production 512 MiB refresh worker must be sampled");
if (values["expect-oom"]) {
assert.match(
errorChain(failure),
/ERR_WORKER_OUT_OF_MEMORY|reaching memory limit|heap out of memory/i,
);
} else {
assert.ifError(failure);
const actual = readRollup(marker);
assert.deepEqual(actual.rollup, reference.rollup);
assert.equal(actual.parsedRecords, reference.parsedRecords);
assert.equal(actual.countedRecords, reference.countedRecords);
assert.equal(actual.checkpoint.kind, "sqlite");
if (actual.checkpoint.kind === "sqlite") {
assert.equal(actual.checkpoint.maxSeq, count);
assert.equal(actual.checkpoint.eventCount, count);
assert.equal(actual.checkpoint.visibleLeafId, `event-${count - 1}`);
}
assert.ok(peakHeapBytes < maxHeapBytes, `Peak ${peakHeapBytes} exceeded ${maxHeapBytes}`);
console.log(
JSON.stringify({
phase: "equality",
rollupEqual: true,
countedRecords: actual.countedRecords,
totalTokens: count * 10,
totalCost: count,
rollupSha256: createHash("sha256").update(JSON.stringify(actual.rollup)).digest("hex"),
reference:
"JSONL worker's existing 128-record aggregation; independently checked token/cost totals",
}),
);
}
} finally {
process.off("worker", observeWorker);
await closeSessionTranscriptReconcileWorkerPool();
await closeOpenClawAgentDatabasesAsync(root);
await closeOpenClawStateDatabaseAsync();
fs.rmSync(root, { recursive: true, force: true });
}