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.
This commit is contained in:
Peter Steinberger 2026-09-23 19:50:28 -07:00 • committed by GitHub
parent e157b5604a
commit 5b90a73498
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
11 changed files with 694 additions and 13 deletions

View file

@ -26,6 +26,12 @@ per-signal state and transport under **Telemetry exporters**. The summary is
redacted and does not include endpoint values, headers, certificates, payloads,
or raw errors.
Doctor reports sessions whose usage-cost cache refresh failed, since their totals
may be incomplete. Check the Gateway logs and request usage again to retry.
The bounded failure history keeps the latest 256 sessions across restarts;
a successful refresh clears that session's warning. `--fix` does not clear a
warning before the session has refreshed successfully.
Related:
- Troubleshooting: [Troubleshooting](/gateway/troubleshooting)

View file

@ -1693,6 +1693,7 @@
"android:version:sync": "node --import ./scripts/tsx.mjs scripts/android-sync-versioning.ts --write",
"audit:seams": "node --import ./scripts/tsx.mjs scripts/audit-seams.mts",
"bench:plugins:invocation": "node --import ./scripts/tsx.mjs scripts/bench-plugin-invocation.ts",
"bench:usage-refresh-memory": "node --import ./scripts/tsx.mjs scripts/bench-usage-refresh-memory.ts",
"build": "node --import ./scripts/tsx.mjs scripts/build-all.mts",
"build:ci-artifacts": "node --import ./scripts/tsx.mjs scripts/build-all.mts ciArtifacts",
"build:docker": "pnpm plugins:assets:build && node --import ./scripts/tsx.mjs scripts/tsdown-build.mts && node --import ./scripts/tsx.mjs scripts/check-cli-bootstrap-imports.mts && node scripts/runtime-postbuild.mjs && node --import ./scripts/tsx.mjs scripts/build-stamp.mts && node --import ./scripts/tsx.mjs scripts/runtime-postbuild-stamp.mts && pnpm plugins:assets:copy && node --import ./scripts/tsx.mjs scripts/write-build-info.ts && node --import ./scripts/tsx.mjs scripts/write-cli-startup-metadata.ts",

View file

@ -0,0 +1,314 @@
// 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 });
}

View file

@ -6,6 +6,7 @@ import { note } from "../../packages/terminal-core/src/note.js";
import { resolveStateDir } from "../config/paths.js";
import { formatErrorMessage, hasErrnoCode } from "../infra/errors.js";
import { deleteSessionCostUsageRollupsExcept } from "../infra/session-cost-usage-cache.sqlite.js";
import { openUsageCostRefreshFailures } from "../infra/session-cost-usage-refresh-health.js";
import { listOpenClawRegisteredAgentDatabases } from "../state/openclaw-agent-db.js";
import { shortenHomePath } from "../utils.js";
import { runDoctorAgentDatabaseOperationAsync } from "./doctor-agent-database-operation.js";
@ -170,6 +171,23 @@ export async function maybeRepairLegacyRuntimeFiles(
shouldRepair: boolean,
env?: NodeJS.ProcessEnv,
): Promise<void> {
const failures = await openUsageCostRefreshFailures(env)
.entries()
.catch((error: unknown) => {
note(
`Could not read usage refresh failure history: ${formatErrorMessage(error)}`,
"Usage cost cache",
);
return [];
});
if (failures.length > 0) {
note(
failures
.map(({ value }) => `- ${value.agentId}: ${value.sessionFile}: ${value.reason}`)
.join("\n"),
"Usage cost cache",
);
}
await maybeScrubConfigAuditLog({ shouldRepair, env });
await maybeRemoveLegacyUsageCostCacheFiles({ shouldRepair, env });
if (shouldRepair) {

View file

@ -5,6 +5,7 @@ import type { Worker } from "node:worker_threads";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { afterEach, expect, it, vi } from "vitest";
import { createDeferred, withTestTimeout } from "../../test/helpers/promise.js";
import { maybeRepairLegacyRuntimeFiles } from "../commands/doctor-usage-cost-cache.js";
import { formatSqliteSessionFileMarker } from "../config/sessions/legacy-sqlite-marker.js";
import {
createSessionEntryWithTranscript,
@ -22,6 +23,7 @@ import { refreshCostUsageCacheForAgent } from "./session-cost-usage-aggregation.
import * as usageCacheSqlite from "./session-cost-usage-cache.sqlite.js";
import { readSessionCostUsageRollupRows } from "./session-cost-usage-cache.test-support.js";
import { resolveUsageCostPricingFingerprint } from "./session-cost-usage-pricing-context.js";
import { openUsageCostRefreshFailures } from "./session-cost-usage-refresh-health.js";
import { prepareUsageCostWorker, runUsageCostWorker } from "./session-cost-usage-worker-runtime.js";
import {
loadCostUsageSummaryFromCache,
@ -31,6 +33,9 @@ import { SqliteWorkerError } from "./sqlite-worker-contract.js";
import { WorkerTaskPool } from "./worker-task-pool.js";
import type { WorkerTaskInput, WorkerTaskOptions } from "./worker-task-pool.types.js";
const note = vi.hoisted(() => vi.fn());
vi.mock("../../packages/terminal-core/src/note.js", () => ({ note }));
const observed = vi.hoisted(() => ({
workers: new Set<Worker>(),
refreshWorkers: new Set<Worker>(),
@ -554,6 +559,47 @@ it("preserves the original host failure when lock cleanup fails and retries that
});
}, 30_000);
it("reports the failed session in doctor and clears it after successful refresh", async () => {
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
const agentId = "usage-failure-health";
const sessionFile = state.path("failed-refresh.jsonl");
await fs.writeFile(sessionFile, usageLine("failed-refresh"));
const prepareLock = usageCacheSqlite.prepareSessionCostUsageRefreshLock;
const observer = vi
.spyOn(usageCacheSqlite, "prepareSessionCostUsageRefreshLock")
.mockImplementation((...args) => ({
...prepareLock(...args),
writeRollup: async () => {
throw new Error("private transcript content");
},
}));
try {
await expect(
refreshCostUsageCacheForAgent({ agentId, sessionFiles: [sessionFile] }),
).rejects.toThrow("private transcript content");
} finally {
observer.mockRestore();
}
const failures = await openUsageCostRefreshFailures(state.env).entries();
expect(failures).toMatchObject([
{ value: { agentId, sessionFile, failedAt: expect.any(Number) } },
]);
expect(JSON.stringify(failures)).not.toContain("private transcript content");
note.mockClear();
await maybeRepairLegacyRuntimeFiles(false, state.env);
expect(note).toHaveBeenCalledWith(expect.stringContaining(sessionFile), "Usage cost cache");
expect(note).toHaveBeenCalledWith(
expect.stringContaining("cached totals may be incomplete"),
"Usage cost cache",
);
await refreshCostUsageCacheForAgent({ agentId, sessionFiles: [sessionFile] });
expect(await openUsageCostRefreshFailures(state.env).entries()).toEqual([]);
note.mockClear();
await maybeRepairLegacyRuntimeFiles(false, state.env);
expect(note.mock.calls.filter(([, title]) => title === "Usage cost cache")).toEqual([]);
});
});
it("retains a late host write failure after cancellation and releases the lock only after settlement", async () => {
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
const agentId = "usage-canceled-write";

View file

@ -0,0 +1,18 @@
import { createCorePluginStateKeyedStore } from "../plugin-state/plugin-state-store.js";
export type UsageCostRefreshFailure = {
agentId: string;
sessionFile: string;
failedAt: number;
reason: string;
};
/** Bounded health facts survive worker/process exits; successful refresh retires each fact. */
export function openUsageCostRefreshFailures(env?: NodeJS.ProcessEnv) {
return createCorePluginStateKeyedStore<UsageCostRefreshFailure>({
ownerId: "core:usage-cost-cache",
namespace: "refresh-failures",
maxEntries: 256,
env,
});
}

View file

@ -0,0 +1,162 @@
import path from "node:path";
import { describe, expect, it } from "vitest";
import { formatSqliteSessionFileMarker } from "../config/sessions/legacy-sqlite-marker.js";
import type { UsageCostRollupEntry } from "./session-cost-usage-rollup-codec.js";
import { scanUsageCostRollupInWorker } from "./session-cost-usage-worker-refresh.js";
const timestamp = Date.parse("2026-09-23T12:00:00.000Z");
type Row = { seq: number; event: Record<string, unknown> };
function message(id: string, parentId: string | null, tokens: number): Row["event"] {
return {
type: "message",
id,
parentId,
timestamp: new Date(timestamp).toISOString(),
message: {
role: "assistant",
usage: { input: tokens, output: 0, totalTokens: tokens, cost: { total: tokens } },
},
};
}
const initial: Row[] = [
{ seq: 1, event: message("root", null, 1) },
{ seq: 4, event: message("a", "root", 2) },
];
async function scan(rows: Row[], previous?: UsageCostRollupEntry) {
const storePath = path.resolve("synthetic-usage-paging.sqlite");
const filePath = formatSqliteSessionFileMarker({
agentId: "main",
sessionId: "paging",
storePath,
});
const maxSeq = rows.at(-1)?.seq ?? 0;
const sizeBytes = rows.reduce(
(sum, row) => sum + Buffer.byteLength(JSON.stringify(row.event)) + 1,
0,
);
return scanUsageCostRollupInWorker({
file: {
kind: "sqlite",
filePath,
sourcePath: filePath,
sessionId: "paging",
maxSeq,
eventCount: rows.length,
size: sizeBytes,
mtimeMs: timestamp,
},
previous,
pricingFingerprint: "synthetic-pricing",
resolveCosts: async (pairs) => pairs.map(() => undefined),
// A short page does not mean end-of-range; decoded byte limits can cut any page.
readRows: async (_marker, afterSeq, throughSeq) =>
rows.filter((row) => row.seq > afterSeq && row.seq <= throughSeq).slice(0, 2),
access: {
readSqliteStats: async () => [
{ maxSeq, eventCount: rows.length, sizeBytes, lastMutationAtMs: timestamp },
],
},
});
}
function expectUsage(entry: UsageCostRollupEntry, tokens: number, records: number, leaf: string) {
expect(entry.parsedRecords).toBe(records);
expect(entry.countedRecords).toBe(records);
expect(
Object.values(entry.rollup.buckets).reduce(
(total, bucket) => total + bucket.totals.totalTokens,
0,
),
).toBe(tokens);
expect(entry.checkpoint).toMatchObject({ kind: "sqlite", visibleLeafId: leaf });
}
describe("paged SQLite usage rollups", () => {
it("selects a branch after reading all short pages and preserves sparse sequence numbers", async () => {
const result = await scan([
{ seq: 1, event: message("root", null, 1) },
{ seq: 4, event: message("hidden", "root", 100) },
{ seq: 5, event: message("hidden-tail", "hidden", 100) },
{ seq: 8, event: { type: "leaf", id: "switch", parentId: "hidden-tail", targetId: "root" } },
{ seq: 13, event: message("chosen", "switch", 2) },
{ seq: 20, event: message("chosen-tail", "chosen", 3) },
]);
expectUsage(result, 6, 3, "chosen-tail");
expect(result.checkpoint).toMatchObject({ maxSeq: 20, eventCount: 6 });
});
it("carries the previous rollup and visible leaf through every append page", async () => {
const previous = await scan(initial);
const result = await scan(
[
...initial,
{ seq: 9, event: message("b", "a", 4) },
{ seq: 12, event: message("c", "b", 8) },
{ seq: 16, event: message("d", "c", 16) },
],
previous,
);
expectUsage(result, 31, 5, "d");
});
it.each([
{
kind: "leaf",
suffix: [{ seq: 20, event: { type: "leaf", id: "switch", parentId: "c", targetId: "root" } }],
tokens: 1,
records: 1,
leaf: "root",
},
{
kind: "reset",
suffix: [
{ seq: 20, event: { type: "reset", id: "reset", parentId: null } },
{ seq: 25, event: message("after-reset", "root", 16) },
],
tokens: 16,
records: 1,
leaf: "after-reset",
},
])(
"rebuilds when a later append page changes the $kind selection",
async ({ suffix, tokens, records, leaf }) => {
const previous = await scan(initial);
const result = await scan(
[
...initial,
{ seq: 9, event: message("b", "a", 4) },
{ seq: 12, event: message("c", "b", 8) },
...suffix,
],
previous,
);
expectUsage(result, tokens, records, leaf);
},
);
it("aggregates duplicate-ID paths in selected ancestry order when sequence numbers go backward", async () => {
const user = (at: number) => ({
type: "message",
id: "user",
parentId: null,
timestamp: new Date(at).toISOString(),
message: { role: "user", content: "synthetic" },
});
const result = await scan([
{ seq: 1, event: user(timestamp) },
{
seq: 2,
event: {
...message("answer", "user", 3),
timestamp: new Date(timestamp + 2000).toISOString(),
},
},
{ seq: 3, event: user(timestamp + 1000) },
{ seq: 4, event: { type: "leaf", id: "switch", parentId: "user", targetId: "answer" } },
]);
expectUsage(result, 3, 1, "answer");
expect(result.rollup.lastUserTimestamp).toBe(timestamp + 1000);
expect(
Object.values(result.rollup.buckets).find((bucket) => bucket.latency.count > 0)?.latency,
).toMatchObject({ count: 1, min: 1000, max: 1000, sum: 1000 });
});
});

View file

@ -364,18 +364,69 @@ async function scanSqliteUsageRollup(params: RollupScanInput): Promise<UsageCost
anchorMatches,
);
const afterSeq = appendCandidate ? (previousCheckpoint?.maxSeq ?? 0) : 0;
const rows = await params.readRows(scope, afterSeq, maxSeq);
const rawRecords = rows.map((row) => row.event).filter(isRecord);
// Branch selection retains only navigation facts, never transcript bodies.
const readNavigation = async (startSeq: number) => {
let after = startSeq;
const records: Array<Record<string, unknown> & { seq: number }> = [];
while (after < maxSeq) {
const page = await params.readRows(scope, after, maxSeq);
if (page.length === 0) {
break;
}
for (const { seq, event } of page) {
const record: Record<string, unknown> & { seq: number } = { seq };
if (isRecord(event)) {
for (const key of [
"type",
"id",
"parentId",
"targetId",
"appendParentId",
"appendMode",
]) {
if (Object.hasOwn(event, key)) {
record[key] = event[key];
}
}
}
records.push(record);
}
after = page.at(-1)!.seq;
}
return records;
};
let navigation = await readNavigation(afterSeq);
const incremental = appendCandidate
? selectIncrementalSqliteRecords(rawRecords, previousCheckpoint?.visibleLeafId)
? selectIncrementalSqliteRecords(navigation, previousCheckpoint?.visibleLeafId)
: undefined;
const appendOnly = Boolean(incremental && params.previous);
const allRows = appendOnly || afterSeq === 0 ? rows : await params.readRows(scope, 0, maxSeq);
const allRecords = appendOnly
? (incremental?.records ?? [])
: selectVisibleTranscriptEvents(allRows.map((row) => row.event)).filter(isRecord);
if (!appendOnly && afterSeq > 0) {
navigation = await readNavigation(0);
}
const selected = appendOnly
? navigation.filter(isCanonicalSessionTranscriptEntry)
: selectVisibleTranscriptEvents(navigation);
const scan = createUsageRollupScan({ ...params, appendOnly });
await scan.addRecords(allRecords);
let page = new Map<number, unknown>();
let pending: Record<string, unknown>[] = [];
for (const record of selected) {
if (!page.has(record.seq)) {
await scan.addRecords(pending);
pending = [];
page.clear();
page = new Map(
(await params.readRows(scope, record.seq - 1, maxSeq)).map((row) => [row.seq, row.event]),
);
if (!page.has(record.seq)) {
throw new Error(`SQLite transcript changed while scanning: ${params.file.filePath}`);
}
}
const event = page.get(record.seq);
if (isRecord(event)) {
pending.push(event);
}
}
await scan.addRecords(pending);
const postFile = await resolveUsageCostTranscriptFile(params.file.filePath, params.access);
if (!postFile || (postFile.maxSeq ?? 0) < maxSeq || (postFile.eventCount ?? 0) < eventCount) {
throw new Error(`SQLite transcript changed while scanning: ${params.file.filePath}`);
@ -389,7 +440,7 @@ async function scanSqliteUsageRollup(params: RollupScanInput): Promise<UsageCost
}
const visibleLeafId = appendOnly
? incremental?.visibleLeafId
: (scanSessionTranscriptTree(allRows.map((row) => row.event)).leafId ?? undefined);
: (scanSessionTranscriptTree(navigation).leafId ?? undefined);
return scan.finish({
kind: "sqlite",
maxSeq,

View file

@ -23,6 +23,7 @@ import { createMemoryTranscriptProjectionSource } from "../config/sessions/sessi
import { withSessionCostUsageWorkerDatabases } from "../config/sessions/session-transcript-worker-runtime.js";
import { resolveStateDir } from "../config/state-dir.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { createSubsystemLogger } from "../logging/subsystem.js";
import { normalizeAgentId } from "../routing/session-key.js";
import { getAsyncWorkSignal } from "../shared/async-work-scope.js";
import type { OpenClawAgentDatabaseOptions } from "../state/openclaw-agent-db-contract.js";
@ -47,6 +48,7 @@ import {
createUsageCostResolver,
resolveUsageCostPricingFingerprint,
} from "./session-cost-usage-pricing-context.js";
import { openUsageCostRefreshFailures } from "./session-cost-usage-refresh-health.js";
import {
UsageCostWorkerReplyError,
type UsageCostWorkerHostEffects,
@ -60,6 +62,7 @@ import type { UsageDailyBucket } from "./session-cost-usage.types.js";
import { withSqliteWorkerCleanupFailure } from "./sqlite-worker-broker-reply.js";
const USAGE_COST_WORKER_TIMEOUT_MS = 5 * 60_000;
const logger = createSubsystemLogger("usage-cost-cache");
export type PreparedUsageCostWorker = {
location: UsageCostWorkerLocation;
@ -328,9 +331,20 @@ export async function runUsageCostWorker(
config: prepared.config,
agentDir: prepared.agentDir,
});
const failures = openUsageCostRefreshFailures(location.env);
const failureKey = (sessionFile: string) =>
JSON.stringify([location.databasePath, sessionFile]);
let activeSessionFile: string | undefined;
const hostErrors = new Map<number, unknown>();
let errorSequence = 0;
try {
const failureEntries = lock
? await failures.entries().catch((error: unknown) => {
logger.warn("Could not read usage refresh failure history", { error });
return [];
})
: [];
const failedKeys = new Set(failureEntries.map((entry) => entry.key));
const result = await scope.run(
{ kind: "usage-cost", location, operation: workerOperation, databases: [] },
{
@ -354,6 +368,13 @@ export async function runUsageCostWorker(
const request = value as UsageCostWorkerHostRequest;
let output: UsageCostWorkerHostEffects[keyof UsageCostWorkerHostEffects]["output"];
switch (request.kind) {
case "refresh-session":
if (!lock) {
throw new Error("Usage report cannot refresh sessions");
}
activeSessionFile = request.input.sessionFile;
output = undefined;
break;
case "pricing":
output = request.input.map(resolveCost);
break;
@ -481,6 +502,16 @@ export async function runUsageCostWorker(
blob: request.input.blob,
updatedAt: request.input.updatedAt,
});
if (output && failedKeys.has(failureKey(request.input.key))) {
await failures
.delete(failureKey(request.input.key), {
assertCurrent: assertRequestCurrent,
})
.catch((error: unknown) => {
logger.warn("Could not clear usage refresh failure fact", { error });
});
}
activeSessionFile = undefined;
break;
default:
throw new Error("Unknown usage worker host request");
@ -503,7 +534,27 @@ export async function runUsageCostWorker(
assertCurrent();
return result;
} catch (error) {
throw restoreWorkerFailure(error, hostErrors);
let failure = restoreWorkerFailure(error, hostErrors);
if (activeSessionFile && !signal?.aborted) {
try {
await failures.register(
failureKey(activeSessionFile),
{
agentId: location.agentId,
sessionFile: activeSessionFile,
failedAt: Date.now(),
reason: "Usage refresh failed; cached totals may be incomplete. Check Gateway logs.",
},
{ assertCurrent },
);
} catch (healthError) {
failure = withSqliteWorkerCleanupFailure(
toErrorObject(failure, "Usage refresh failed"),
healthError,
);
}
}
throw failure;
}
});
}

View file

@ -21,7 +21,7 @@ import {
} from "../state/openclaw-agent-db-readonly.js";
import { isIncognitoOpenClawAgentSqlitePath } from "../state/openclaw-agent-db.js";
import { encodeOpenClawStateWorkerError } from "../state/openclaw-state-worker-error.js";
import { executeSqliteQuerySync } from "./kysely-sync.js";
import { iterateSqliteQuerySync } from "./kysely-sync.js";
import {
readSessionCostUsageRollupRowsInDatabase,
readSessionCostUsageRollupBodyInDatabase,
@ -493,8 +493,20 @@ export async function executeUsageCostWorker(
.where("session_id", "=", marker.sessionId)
.where("seq", ">", afterSeq)
.where("seq", "<=", throughSeq)
.orderBy("seq", "asc");
return executeSqliteQuerySync(opened.db, query).rows;
.orderBy("seq", "asc")
.limit(1_024);
const page: Array<{ seq: number; event_json: string }> = [];
let bytes = 0;
// Stop before parsing: retain at most 8 MiB plus one lookahead event.
for (const row of iterateSqliteQuerySync(opened.db, query)) {
const size = Buffer.byteLength(row.event_json);
if (page.length > 0 && bytes + size > 8 * 1024 * 1024) {
break;
}
page.push(row);
bytes += size;
}
return page;
}),
{ ...database, env },
);
@ -514,6 +526,7 @@ export async function executeUsageCostWorker(
};
for (const { file, row, envelope, rebuild } of stale.slice(0, maxFiles)) {
control.throwIfCancelled();
await host("refresh-session", { sessionFile: file.filePath });
let previous: UsageCostRollupEntry | undefined;
if (
!rebuild &&

View file

@ -108,6 +108,7 @@ type UsageCostPruneRow = {
};
export type UsageCostWorkerHostEffects = {
"refresh-session": { input: { sessionFile: string }; output: void };
pricing: {
input: Array<{ provider?: string; model?: string }>;
output: Array<ModelCostConfig | undefined>;