fix(sessions): order WAL recovery and join probe release

Keep checkpoint completion ordering in the WAL owner using a private
monotonic snapshot relayed unchanged between threads. Public health and
disk-budget diagnostics retain epoch timestamps and JSON-safe fields.
This lets successful checkpoints release cleanup deferral after a system
clock correction while rejecting delayed older completion facts.

Join the parent's real commit-settlement probe release before worker
post-commit page maintenance. The existing shared gate distinguishes child
transaction settlement from parent writer release, preserves terminal
acknowledgments, and reports failed release without claiming success.
The native writer-lock proof and original approval deadline remain intact.

Controlled real-SQLite negative cases reproduce both defects independently.
The candidate passes owner and adjacent cleanup, cold-storage, writer,
and public-status tests. Historical CI scheduling remains unobserved;
these controls establish the repaired invariants, not historical attribution.

Final main-based Testbox proof: 61 tests; core and affected test types;
formatting; targeted lint. Broader invariant proof also passed 209 adjacent
tests and all 25 core-test type graphs. The only full-gate failure was the
now-corrected no-unsafe-finally error path, independently reviewed and
revalidated without suppressions.

Standalone pnpm test --maxWorkers=1 command wall on Node 24.19.0, 4 CPUs,
about 16 GB RAM, fresh lease with a cold dedicated module cache warming in
this order:
- sqlite-wal-checkpoint.test.ts: 3.480 s
- sqlite-wal-reclamation.test.ts: 2.286 s
- session-accessor.sqlite-reclamation-commit.test.ts: 12.492 s
- session-accessor.sqlite-reclamation-wal.test.ts: 25.158 s
- session-history-wal-budget.test.ts: 24.058 s

Preserve the retired-worker cleanup and worker-backed fixture changes from
#154889 and #155028. Final integration rechecks reclamation reuse and exact
lease recovery alongside all five WAL owner files.
This commit is contained in:
Peter Steinberger 2026-09-21 12:06:23 -07:00
parent aca0c60535
commit bef4401fbc
20 changed files with 444 additions and 93 deletions

View file

@ -344,7 +344,8 @@ The result records `deferredReason: "checkpoint-incomplete"`, WAL bytes before a
after, and the checkpoint outcome. Automatic and manual budget passes remain
deferred until the checkpoint owner observes a completed checkpoint; elapsed time
or a budget change alone does not retry pruning. Normal periodic checkpointing
continues, and subsequent activity can resume cleanup after recovery.
continues, and subsequent activity can resume cleanup after recovery, including
after a system clock correction.
Look for `session history disk budget deferred until a completed WAL checkpoint is observed`
in the Gateway log. Its checkpoint fields include bounded operation names for

View file

@ -1449,6 +1449,12 @@ completion, checkpoint facts, and physical bytes before and after. Budget cleanu
remains deferred until the checkpoint owner reports completion, preserving retained
data instead of adding writes behind a pinned WAL.
Checkpoint ordering uses a private monotonic observation shared by the host and
its workers; health timestamps remain wall-clock diagnostics. Post-commit page
maintenance also waits for the parent's commit-settlement probe to release its
writer lock. Child transaction settlement and parent probe release are distinct
facts in the existing commit gate; failed release cannot acknowledge success.
Queued archive pruning prepares cold connections through the same asynchronous
admission owner while retaining its existing writer section. File-backed page
drains use the existing reclamation worker, acquired before the caller's writer

View file

@ -33,6 +33,7 @@ export type ReclamationDatabaseOptions = OpenClawAgentDatabaseOptions & {
export type SqliteSessionReclamationCallbacks = {
beforeMutation?: () => void;
onCommit?: (database: OpenClawAgentDatabase, result?: SqliteSessionReclamationResult) => void;
afterCommit?: () => void;
};
export type ReclamationDeleteParams = Omit<DeleteSessionEntryLifecycleParams, "commitGuard">;

View file

@ -34,6 +34,7 @@ import type { SqliteSessionReclamationPlan } from "./session-accessor.sqlite-lif
import {
markSqliteReclamationSettled,
waitForSqliteReclamationCommit,
waitForSqliteReclamationParentRelease,
} from "./session-accessor.sqlite-reclamation-commit.js";
import type {
SqliteCanonicalValidationWorkerRequest,
@ -130,8 +131,7 @@ async function runColdMutationWorker(port: MessagePort, data: SessionColdWorkerD
);
},
);
// The parent joins the cold transaction, not the subsequent bounded page drain.
markSqliteReclamationSettled(commitGate);
waitForSqliteReclamationParentRelease(commitGate);
if (data.plan.kind !== "cold-restore") {
await reclaimSqliteFreePages(data.plan.databaseOptions, undefined, { maxPasses: 64 });
}
@ -189,12 +189,12 @@ export async function runReclamationWorkerPort(
let failureCleanup: Awaited<ReturnType<typeof settleReclamationDatabase>> | undefined;
let checkpointResultOwnedByRequest = false;
const checkpointPath = sqliteReaderDatabasePathKey(databaseOptions.path);
const stopCheckpointRelay = onSqliteWalCheckpoint(({ databasePath, health }) => {
const stopCheckpointRelay = onSqliteWalCheckpoint(({ databasePath, ...snapshot }) => {
if (!checkpointResultOwnedByRequest && databasePath === checkpointPath && claim?.isCurrent()) {
port.postMessage({
type: "checkpoint",
operationId,
health,
snapshot,
} satisfies SqliteReclamationWorkerMessage);
}
});
@ -378,6 +378,8 @@ export async function runReclamationWorkerPort(
{
beforeMutation: currentClaim.assertCurrent,
onCommit: authorizeCommit,
afterCommit: () =>
waitForSqliteReclamationParentRelease(request.commitGate),
},
);
// Warm results must not revive proof invalidated by the parent between requests.

View file

@ -1,4 +1,4 @@
import { publishSqliteWalCheckpointHealth } from "../../infra/sqlite-wal-checkpoint.js";
import { publishSqliteWalCheckpointObservation } from "../../infra/sqlite-wal-checkpoint.js";
import type { SqliteWalReclamationResult } from "../../infra/sqlite-wal.js";
import { readOpenClawAgentDatabaseIdentity } from "../../state/openclaw-agent-db-identity.js";
import { retainOpenClawAgentDatabaseReadOnly } from "../../state/openclaw-agent-db-readonly.js";
@ -94,7 +94,7 @@ export async function withSqliteSessionPageReclamation<T>(
throw new Error("SQLite page reclamation returned another operation's result");
}
if (result.value.checkpoint) {
result.value.checkpoint = publishSqliteWalCheckpointHealth(
result.value.checkpoint = publishSqliteWalCheckpointObservation(
databaseOptions.path,
result.value.checkpoint,
);

View file

@ -1,8 +1,10 @@
import { parentPort, workerData } from "node:worker_threads";
import { openNodeSqliteDatabase } from "../../infra/node-sqlite.js";
import { configureSqliteWalMaintenance } from "../../infra/sqlite-wal.js";
import {
markSqliteReclamationSettled,
waitForSqliteReclamationCommit,
waitForSqliteReclamationParentRelease,
} from "./session-accessor.sqlite-reclamation-commit.js";
type CommitFixture = {
@ -10,6 +12,7 @@ type CommitFixture = {
gate: SharedArrayBuffer;
progress: SharedArrayBuffer;
holdAfterApproval?: boolean;
checkpointAfterCommit?: boolean;
outcome?: "rollback" | "exit-before-commit" | "exit-after-commit";
};
@ -21,6 +24,13 @@ if (!port) {
const fixture = workerData as CommitFixture;
const progress = new Int32Array(fixture.progress);
const database = openNodeSqliteDatabase(fixture.databasePath);
const maintenance = fixture.checkpointAfterCommit
? configureSqliteWalMaintenance(database, {
databasePath: fixture.databasePath,
checkpointIntervalMs: 0,
busyTimeoutMs: 0,
})
: undefined;
database.exec("BEGIN IMMEDIATE; UPDATE proof SET value = 2");
try {
waitForSqliteReclamationCommit(fixture.gate, () => port.postMessage("commit-request"));
@ -39,12 +49,25 @@ try {
if (fixture.outcome === "exit-after-commit") {
process.exit(9);
}
if (maintenance) {
Atomics.wait(progress, 1, 0);
waitForSqliteReclamationParentRelease(fixture.gate);
port.postMessage({
checkpointCompleted: maintenance.checkpoint(),
health: maintenance.health,
parentRelease: Atomics.load(new Int32Array(fixture.gate), 0),
});
}
} catch (error) {
if (database.isTransaction) {
database.exec("ROLLBACK");
}
port.postMessage({ error: String(error) });
port.postMessage({
error: String(error),
parentRelease: Atomics.load(new Int32Array(fixture.gate), 0),
});
} finally {
maintenance?.close();
database.close();
markSqliteReclamationSettled(fixture.gate);
Atomics.store(progress, 0, 2);

View file

@ -7,7 +7,11 @@ import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"
import { openNodeSqliteDatabase } from "../../infra/node-sqlite.js";
import * as nodeSqlite from "../../infra/node-sqlite.js";
import * as sqliteTransaction from "../../infra/sqlite-transaction.js";
import { withSqliteReclamationAuthorization } from "./session-accessor.sqlite-reclamation-commit.js";
import {
revokeSqliteReclamationCommit,
waitForSqliteReclamationParentRelease,
withSqliteReclamationAuthorization,
} from "./session-accessor.sqlite-reclamation-commit.js";
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
afterEach(() => vi.restoreAllMocks());
@ -26,6 +30,7 @@ function waitForApproval(progress: Int32Array): void {
function createCommitFixture(
options: {
holdAfterApproval?: boolean;
checkpointAfterCommit?: boolean;
outcome?: "rollback" | "exit-before-commit" | "exit-after-commit";
} = {},
) {
@ -75,6 +80,106 @@ function createCommitFixture(
};
}
test("does not wait for a parent acknowledgment without a consumed commit approval", () => {
const gate = new SharedArrayBuffer(Int32Array.BYTES_PER_ELEMENT);
waitForSqliteReclamationParentRelease(gate);
expect(Atomics.load(new Int32Array(gate), 0)).toBe(0);
revokeSqliteReclamationCommit(gate);
expect(() => waitForSqliteReclamationParentRelease(gate)).toThrow("commit was not authorized");
});
test.each([false, true])(
"joins parent probe release before checkpoint (release fails: %s)",
async (releaseFails) => {
const fixture = createCommitFixture({ checkpointAfterCommit: true });
const transact = sqliteTransaction.runSqliteImmediateTransactionSync;
let heldWriter = false;
let releaseProbe: (() => void) | undefined;
let probeStillHeld: (() => boolean) | undefined;
try {
await fixture.requested;
const checkpoint = once(fixture.worker, "message");
if (releaseFails) {
const open = nodeSqlite.openNodeSqliteDatabase;
vi.spyOn(nodeSqlite, "openNodeSqliteDatabase").mockImplementationOnce((...args) => {
const probe = open(...args);
const exec = probe.exec.bind(probe);
const close = probe.close.bind(probe);
probeStillHeld = () => probe.isOpen && probe.isTransaction;
releaseProbe = () => {
if (probe.isOpen) {
if (probe.isTransaction) {
exec("ROLLBACK");
}
close();
}
};
vi.spyOn(probe, "exec").mockImplementation((sql) => {
if (sql === "COMMIT" || sql === "ROLLBACK") {
throw new Error(`injected probe ${sql} failure`);
}
return exec(sql);
});
vi.spyOn(probe, "close").mockImplementation(() => {
throw new Error("injected probe close failure");
});
return probe;
});
}
vi.spyOn(sqliteTransaction, "runSqliteImmediateTransactionSync").mockImplementation(
(db, operation, options) =>
transact(
db,
() => {
const value = operation();
heldWriter = db.isTransaction;
const gate = new Int32Array(fixture.gate);
const committing = Atomics.load(gate, 0);
fixture.release();
Atomics.wait(gate, 0, committing, 5_000);
expect(Atomics.load(gate, 0)).not.toBe(committing);
return value;
},
options,
),
);
await fixture.withAuthorization(
() => {},
async (authorize) => {
const errors = authorize();
if (releaseFails) {
expect(errors.map(String)).toEqual(
expect.arrayContaining([
"Error: injected probe COMMIT failure",
"Error: injected probe close failure",
"Error: SQLite commit-settlement probe remains in a transaction after close",
]),
);
expect(probeStillHeld?.()).toBe(true);
} else {
expect(errors).toEqual([]);
}
},
);
const [observed] = await checkpoint;
expect(heldWriter).toBe(true);
expect(observed).toMatchObject(
releaseFails
? {
error: "Error: SQLite parent commit-settlement probe did not release its writer lock",
}
: { checkpointCompleted: true, health: { state: "complete", walBytes: 0 } },
);
expect(await fixture.exited).toEqual([0]);
expect(Atomics.load(new Int32Array(fixture.gate), 0)).toBe(observed.parentRelease);
expect(fixture.value()).toBe(2);
} finally {
releaseProbe?.();
await fixture.close();
}
},
);
test.each(["transient", "permanent", "closed"] as const)(
"joins a consumed commit despite a %s barrier failure",
async (fault) => {

View file

@ -15,6 +15,8 @@ const REJECTED = 2;
const COMMITTING = 3;
const SETTLED = 4;
const REQUESTED = 5;
const PARENT_RELEASED = 6;
const PARENT_RELEASE_FAILED = 7;
/** Preserve the reclamation owner's context when an unrelated synchronous writer helps. */
export async function withSqliteReclamationAuthorization<T>(
@ -98,8 +100,44 @@ export function waitForSqliteReclamationCommit(
export function markSqliteReclamationSettled(buffer: SharedArrayBuffer | undefined): void {
if (buffer) {
const shared = new Int32Array(buffer);
Atomics.store(shared, 0, SETTLED);
Atomics.notify(shared, 0);
let current = Atomics.load(shared, 0);
while (current !== PARENT_RELEASED && current !== PARENT_RELEASE_FAILED) {
const observed = Atomics.compareExchange(shared, 0, current, SETTLED);
if (observed === current) {
Atomics.notify(shared, 0);
return;
}
current = observed;
}
}
}
/** Post-commit maintenance must not contend with the parent's own settlement probe. */
export function waitForSqliteReclamationParentRelease(buffer: SharedArrayBuffer): void {
const shared = new Int32Array(buffer);
const decision = Atomics.load(shared, 0);
// Nonmutating requests never ask the parent for commit approval.
if (decision === WAITING) {
return;
}
if (
decision !== COMMITTING &&
decision !== SETTLED &&
decision !== PARENT_RELEASED &&
decision !== PARENT_RELEASE_FAILED
) {
throw new Error("SQLite session reclamation commit was not authorized");
}
markSqliteReclamationSettled(buffer);
for (;;) {
const phase = Atomics.load(shared, 0);
if (phase === PARENT_RELEASED) {
return;
}
if (phase === PARENT_RELEASE_FAILED) {
throw new Error("SQLite parent commit-settlement probe did not release its writer lock");
}
Atomics.wait(shared, 0, SETTLED);
}
}
@ -115,6 +153,7 @@ function authorizeSqliteReclamationCommit(
let database: DatabaseSync | undefined;
const recoveredErrors: unknown[] = [];
let settled = false;
let failure: { error: unknown } | undefined;
try {
database = openNodeSqliteDatabase(databasePath);
setSqliteBusyTimeout(database, COMMIT_DECISION_TIMEOUT_MS);
@ -165,20 +204,40 @@ function authorizeSqliteReclamationCommit(
}
}
} catch (error) {
failure = { error };
rejectCommit(shared);
throw error;
} finally {
let closeFailure: { error: unknown } | undefined;
try {
if (database?.isOpen) {
database.close();
}
} catch (error) {
// The original authorization failure stays fatal. After settlement, the
// Worker's result owns success and all postcommit publication must continue.
if (settled) {
recoveredErrors.push(error);
}
closeFailure = { error };
}
if (settled) {
// Child settlement alone does not release a probe whose own transaction failed.
const released = !database?.isOpen || !database.isTransaction;
Atomics.store(shared, 0, released ? PARENT_RELEASED : PARENT_RELEASE_FAILED);
Atomics.notify(shared, 0);
if (closeFailure) {
recoveredErrors.push(closeFailure.error);
}
if (!released) {
recoveredErrors.push(
new Error("SQLite commit-settlement probe remains in a transaction after close"),
);
}
} else if (failure && closeFailure) {
failure = {
error: new AggregateError([failure.error, closeFailure.error], String(failure.error), {
cause: failure.error,
}),
};
}
}
if (failure) {
throw failure.error;
}
return recoveredErrors;
}

View file

@ -5,7 +5,9 @@ import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { afterEach, expect, test, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js";
import { requireNodeSqlite } from "../../infra/node-sqlite.js";
import { onSqliteWalCheckpoint } from "../../infra/sqlite-wal-checkpoint.js";
import * as nodeSqlite from "../../infra/node-sqlite.js";
import { sqliteReaderDatabasePathKey } from "../../infra/sqlite-reader-lifecycle.js";
import * as walCheckpoint from "../../infra/sqlite-wal-checkpoint.js";
import { configureSqliteWalMaintenance } from "../../infra/sqlite-wal.js";
import { closeCachedOpenClawAgentDatabase } from "../../state/openclaw-agent-db-lifecycle.js";
import { withOpenClawAgentDatabaseReadOnly } from "../../state/openclaw-agent-db-readonly.js";
@ -31,6 +33,7 @@ import {
const hooks = vi.hoisted(() => ({
beforeAuthorization: undefined as (() => void) | undefined,
afterAuthorization: undefined as (() => void) | undefined,
onWorker: undefined as ((worker: SqliteReclamationWorker) => void) | undefined,
}));
vi.mock("./session-accessor.sqlite-reclamation-worker.js", async (importOriginal) => {
@ -50,7 +53,11 @@ vi.mock("./session-accessor.sqlite-reclamation-worker.js", async (importOriginal
...params,
onCommitRequest: () => {
hooks.beforeAuthorization?.();
return params.onCommitRequest();
try {
return params.onCommitRequest();
} finally {
hooks.afterAuthorization?.();
}
},
}),
);
@ -68,6 +75,7 @@ vi.mock("./session-accessor.sqlite-reclamation-worker.js", async (importOriginal
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
afterEach(async () => {
hooks.beforeAuthorization = undefined;
hooks.afterAuthorization = undefined;
hooks.onWorker = undefined;
await closeOpenClawAgentDatabasesAsync();
closeOpenClawAgentDatabasesForTest();
@ -91,6 +99,7 @@ test.each(["reclaim", "worker-close"] as const)(
"reclaims pages off-thread and releases budget deferral through %s",
async (recovery) => {
const { database, databaseOptions } = createFixture();
const databasePathKey = sqliteReaderDatabasePathKey(database.path);
database.db.exec(`INSERT INTO cache_entries(scope, key, blob, updated_at)
VALUES ('wal-proof', 'pages', zeroblob(4194304), 1);
DELETE FROM cache_entries WHERE scope = 'wal-proof';`);
@ -107,6 +116,51 @@ test.each(["reclaim", "worker-close"] as const)(
maintenance: { maxDiskBytes: 1, highWaterBytes: 1 },
};
const budget = getBudgetKickState(params.storePath, params.maintenance);
let capturingProbe = false;
let parentReleaseNs: bigint | undefined;
hooks.beforeAuthorization = () => {
capturingProbe = true;
};
hooks.afterAuthorization = () => {
capturingProbe = false;
};
const open = nodeSqlite.openNodeSqliteDatabase;
const observeProbe = vi
.spyOn(nodeSqlite, "openNodeSqliteDatabase")
.mockImplementation((...args) => {
const opened = open(...args);
if (capturingProbe && sqliteReaderDatabasePathKey(args[0]) === databasePathKey) {
const close = opened.close.bind(opened);
vi.spyOn(opened, "close").mockImplementation(() => {
close();
parentReleaseNs = process.hrtime.bigint();
});
}
return opened;
});
let completedAt: number | undefined;
const relayedCompletions: number[] = [];
const publish = walCheckpoint.publishSqliteWalCheckpointObservation;
const relay = vi
.spyOn(walCheckpoint, "publishSqliteWalCheckpointObservation")
.mockImplementation((databasePath, snapshot) => {
if (
sqliteReaderDatabasePathKey(databasePath) !== databasePathKey ||
snapshot.health.state !== "complete" ||
completedAt === undefined
) {
return publish(databasePath, snapshot);
}
relayedCompletions.push(completedAt);
return publish(databasePath, {
...snapshot,
health: {
...snapshot.health,
observedAtMs: completedAt,
lastCompletedAtMs: completedAt,
},
});
});
let retainedWorker: SqliteReclamationWorker | undefined;
hooks.onWorker = (worker) => {
retainedWorker = worker;
@ -133,16 +187,21 @@ test.each(["reclaim", "worker-close"] as const)(
const blocked = await reclaim();
expect(blocked).toMatchObject({
checkpointCompleted: false,
checkpoint: { state: "blocked" },
checkpoint: { health: { state: "blocked" } },
checkpointIncomplete: 1,
vacuumPasses: 0,
});
expect(freePages()).toBe(original);
deferPhysicalBudgetForCheckpoint(params, database.path, blocked.checkpoint);
assert(blocked.checkpoint);
completedAt = blocked.checkpoint.health.observedAtMs - 1;
if (recovery === "reclaim") {
reader.exec("ROLLBACK");
const completed = await reclaim(7);
expect(completed.vacuumPagesRequested).toBe(7);
assert(parentReleaseNs !== undefined);
assert(completed.checkpoint);
expect(completed.checkpoint.observedAtNs).toBeGreaterThanOrEqual(parentReleaseNs);
expect(original - freePages()).toBeGreaterThan(0);
expect(original - freePages()).toBeLessThanOrEqual(7);
}
@ -161,8 +220,8 @@ test.each(["reclaim", "worker-close"] as const)(
expect(database.db.isOpen).toBe(false);
expect(budget.checkpointBlocked).toBeDefined();
const observed: string[] = [];
const unsubscribe = onSqliteWalCheckpoint(({ databasePath, health }) => {
if (databasePath === database.path) {
const unsubscribe = walCheckpoint.onSqliteWalCheckpoint(({ databasePath, health }) => {
if (databasePath === databasePathKey) {
observed.push(health.state);
}
});
@ -175,14 +234,17 @@ test.each(["reclaim", "worker-close"] as const)(
} finally {
unsubscribe();
}
expect(budget.checkpointBlocked).toBeUndefined();
}
expect(relayedCompletions.length).toBeGreaterThan(0);
expect(budget.checkpointBlocked).toBeUndefined();
} finally {
if (reader.isTransaction) {
reader.exec("ROLLBACK");
}
reader.close();
parentReclaim.mockRestore();
relay.mockRestore();
observeProbe.mockRestore();
}
},
);

View file

@ -5,8 +5,8 @@ import { isDeepStrictEqual } from "node:util";
import { toStringifiedError } from "@openclaw/normalization-core/error-coercion";
import { SQLITE_IDLE_HANDLE_TTL_MS } from "../../infra/sqlite-handle-lifecycle.js";
import {
publishSqliteWalCheckpointHealth,
type SqliteWalHealth,
publishSqliteWalCheckpointObservation,
type SqliteWalCheckpointSnapshot,
} from "../../infra/sqlite-wal-checkpoint.js";
import { captureStateDatabaseCoordinatorRuntime } from "../../infra/state-database-coordinator.js";
import { createSubsystemLogger } from "../../logging/subsystem.js";
@ -95,7 +95,7 @@ type WorkerCleanup = { cleanupWarnings: string[]; settled: boolean };
export type SqliteReclamationWorkerMessage =
| SqliteMutationWorkerMessage<SqliteSessionReclamationResult>
| { type: "lease"; receipt: OpenClawAgentDatabaseWorkerLeaseReceipt }
| { type: "checkpoint"; operationId: number; health: SqliteWalHealth }
| { type: "checkpoint"; operationId: number; snapshot: SqliteWalCheckpointSnapshot }
| ({ type: "closed" } & WorkerCleanup);
const log = createSubsystemLogger("session-sqlite");
@ -508,7 +508,7 @@ export class SqliteReclamationWorker {
this.stateContext.admission.assertCurrent();
this.assertPathCurrent();
// Native close keeps its admitted custody after new requests are revoked.
publishSqliteWalCheckpointHealth(this.options.path, message.health);
publishSqliteWalCheckpointObservation(this.options.path, message.snapshot);
} catch {
// A retired state owner cannot publish late diagnostics or change native settlement.
}

View file

@ -267,9 +267,22 @@ export function reclaimSqliteSessionInTransaction(
maxPages: plan.maxPages,
beforeMutation: callbacks.beforeMutation,
onCommit: () => callbacks.onCommit?.(database),
afterCommit: callbacks.afterCommit,
}),
};
}
const result = reclaimSqliteRowsInTransaction(plan, callbacks);
callbacks.afterCommit?.();
if (result.kind === "history-eviction" && result.value.deleted) {
reclaimSqliteFreePagesBestEffort(plan.databaseOptions);
}
return result;
}
function reclaimSqliteRowsInTransaction(
plan: Exclude<SqliteSessionReclamationPlan, { kind: "maintenance-pages" }>,
callbacks: SqliteSessionReclamationCallbacks,
): SqliteSessionReclamationResult {
if (
plan.kind === "maintenance-plan" ||
plan.kind === "maintenance-finalize" ||
@ -391,9 +404,6 @@ export function reclaimSqliteSessionInTransaction(
}
return { archivedTranscripts: deleted ? archivedTranscripts : [], deleted };
}, plan.databaseOptions);
if (plan.kind === "history-eviction" && value.deleted) {
reclaimSqliteFreePagesBestEffort(plan.databaseOptions);
}
return { kind: plan.kind, value };
}

View file

@ -2,7 +2,10 @@ import fs from "node:fs";
import path from "node:path";
import { setImmediate } from "node:timers/promises";
import { executeSqliteQuerySync } from "../../infra/kysely-sync.js";
import type { SqliteWalHealth } from "../../infra/sqlite-wal-checkpoint.js";
import type {
SqliteWalCheckpointSnapshot,
SqliteWalHealth,
} from "../../infra/sqlite-wal-checkpoint.js";
import type { SqliteWalReclamationResult } from "../../infra/sqlite-wal-reclamation.js";
import {
openOpenClawAgentDatabase,
@ -49,7 +52,7 @@ async function withArchivePruningDatabase<T>(
type PageReclamation = {
reclaimPages?: (maxPages?: number) => Promise<SqliteWalReclamationResult>;
onCheckpointIncomplete?: (checkpoint: SqliteWalHealth | undefined) => void;
onCheckpointIncomplete?: (checkpoint: SqliteWalCheckpointSnapshot | undefined) => void;
};
export type SessionArchivePruningResult = {
@ -94,7 +97,7 @@ export async function reclaimSqliteFreePages(
diagnostics.checkpointMaxMs ?? 0,
result.checkpointMaxMs,
);
diagnostics.checkpoint = result.checkpoint;
diagnostics.checkpoint = result.checkpoint?.health;
}
if (!result.checkpointCompleted) {
limits?.onCheckpointIncomplete?.(result.checkpoint);

View file

@ -1,5 +1,9 @@
import { sqliteReaderDatabasePathKey } from "../../infra/sqlite-reader-lifecycle.js";
import { onSqliteWalCheckpoint, type SqliteWalHealth } from "../../infra/sqlite-wal-checkpoint.js";
import {
onSqliteWalCheckpoint,
type SqliteWalCheckpointSnapshot,
type SqliteWalHealth,
} from "../../infra/sqlite-wal-checkpoint.js";
import { createSubsystemLogger } from "../../logging/subsystem.js";
import type { SessionDiskBudgetSweepResult } from "./disk-budget.types.js";
import type { ResolvedSessionMaintenanceConfig } from "./store-maintenance.js";
@ -66,12 +70,16 @@ type BudgetKickState = {
running: boolean;
pendingForce?: SessionHistoryBudgetKick;
blockedUntil?: number;
checkpointBlocked?: { databasePath: string; checkpoint?: SqliteWalHealth; warned?: boolean };
checkpointBlocked?: {
databasePath: string;
checkpoint?: SqliteWalCheckpointSnapshot;
warned?: boolean;
};
};
export const budgetKickStateByStore = new Map<string, BudgetKickState>();
onSqliteWalCheckpoint(({ databasePath, health }) => {
onSqliteWalCheckpoint(({ databasePath, health, observedAtNs }) => {
if (health.state !== "complete") {
return;
}
@ -79,7 +87,8 @@ onSqliteWalCheckpoint(({ databasePath, health }) => {
// Worker messages can arrive after a newer parent-side observation.
if (
state.checkpointBlocked?.databasePath === databasePath &&
health.observedAtMs >= (state.checkpointBlocked.checkpoint?.observedAtMs ?? -Infinity)
(!state.checkpointBlocked.checkpoint ||
observedAtNs >= state.checkpointBlocked.checkpoint.observedAtNs)
) {
state.checkpointBlocked = undefined;
state.blockedUntil = undefined;
@ -92,7 +101,7 @@ onSqliteWalCheckpoint(({ databasePath, health }) => {
export function deferPhysicalBudgetForCheckpoint(
params: SessionHistoryDiskBudgetParams,
databasePath: string,
checkpoint: SqliteWalHealth | undefined,
checkpoint: SqliteWalCheckpointSnapshot | undefined,
): void {
const state = getBudgetKickState(params.storePath, params.maintenance);
state.checkpointBlocked = { databasePath: sqliteReaderDatabasePathKey(databasePath), checkpoint };

View file

@ -101,7 +101,7 @@ export async function inspectSqliteSessionHistoryDiskBudget(
diskBudget: {
...diskBudget,
deferredReason: "checkpoint-incomplete",
checkpoint: blocked.checkpoint,
checkpoint: blocked.checkpoint?.health,
walBytesBefore: usage.databaseWalBytes,
walBytesAfter: usage.databaseWalBytes,
},
@ -392,7 +392,7 @@ async function enforceSessionHistoryMaintenanceSerialized(
maxBytes: maxDiskBytes,
highWaterBytes,
deferred: {
checkpoint: blocked.checkpoint,
checkpoint: blocked.checkpoint?.health,
walBytesBefore: initialUsage.databaseWalBytes,
walBytesAfter: initialUsage.databaseWalBytes,
},

View file

@ -2,11 +2,19 @@ import assert from "node:assert/strict";
import { channel } from "node:diagnostics_channel";
import fs from "node:fs";
import path from "node:path";
import { performance } from "node:perf_hooks";
import { afterEach, expect, it, vi } from "vitest";
import { getNodeSqliteKysely, iterateSqliteQuerySync } from "../../infra/kysely-sync.js";
import { openNodeSqliteDatabase } from "../../infra/node-sqlite.js";
import { withSqliteReaderOwner } from "../../infra/sqlite-reader-lifecycle.js";
import { publishSqliteWalCheckpointHealth } from "../../infra/sqlite-wal-checkpoint.js";
import {
sqliteReaderDatabasePathKey,
withSqliteReaderOwner,
} from "../../infra/sqlite-reader-lifecycle.js";
import {
onSqliteWalCheckpoint,
publishSqliteWalCheckpointObservation,
type SqliteWalCheckpointSnapshot,
} from "../../infra/sqlite-wal-checkpoint.js";
import {
closeOpenClawAgentDatabasesAsync,
openOpenClawAgentDatabase,
@ -58,6 +66,7 @@ it.each(["transaction", "iterator"] as const)(
});
const options = { agentId: "main", env: state.env };
const database = openOpenClawAgentDatabase(options);
const databasePathKey = sqliteReaderDatabasePathKey(database.path);
ensureSessionTranscriptArchiveSchema(database.db);
database.db.exec("PRAGMA wal_autocheckpoint=0");
const sessions = state.sessionsDir();
@ -73,8 +82,18 @@ it.each(["transaction", "iterator"] as const)(
fs.writeFileSync(path.join(sessions, name), bytes);
return name;
});
expect(database.walMaintenance.checkpoint()).toBe(true);
const previouslyCompleted = database.walMaintenance.health;
const initialCheckpoints: SqliteWalCheckpointSnapshot[] = [];
const stopObserving = onSqliteWalCheckpoint(({ databasePath, health, observedAtNs }) => {
if (databasePath === databasePathKey) {
initialCheckpoints.push({ health, observedAtNs });
}
});
try {
expect(database.walMaintenance.checkpoint()).toBe(true);
} finally {
stopObserving();
}
const previouslyCompleted = initialCheckpoints.at(-1);
assert(previouslyCompleted);
const initial = await measureSessionPhysicalDiskUsage(storePath);
const maintenance = resolveMaintenanceConfigFromInput({
@ -179,10 +198,14 @@ it.each(["transaction", "iterator"] as const)(
}),
]);
assert(blocked?.checkpoint);
// A delayed worker fact must not clear a newer incomplete observation.
publishSqliteWalCheckpointHealth(database.path, {
// A delayed worker fact remains older even when its wall clock was ahead.
publishSqliteWalCheckpointObservation(database.path, {
...previouslyCompleted,
observedAtMs: blocked.checkpoint.observedAtMs - 1,
health: {
...previouslyCompleted.health,
observedAtMs: blocked.checkpoint.observedAtMs + 3600_000,
lastCompletedAtMs: blocked.checkpoint.observedAtMs + 3600_000,
},
});
for (let hour = 1; hour <= 3; hour++) {
kickSessionHistoryDiskBudgetMaintenance({
@ -224,7 +247,14 @@ it.each(["transaction", "iterator"] as const)(
release();
// Release alone is not an observed successful checkpoint.
await expect(enforce()).resolves.toMatchObject({ deferredReason: "checkpoint-incomplete" });
expect(database.walMaintenance.checkpoint()).toBe(true);
const checkpointClock = vi
.spyOn(Date, "now")
.mockReturnValue(blocked.checkpoint.observedAtMs - 1);
try {
expect(database.walMaintenance.checkpoint()).toBe(true);
} finally {
checkpointClock.mockRestore();
}
expect(fs.statSync(`${database.path}-wal`).size).toBe(0);
const recovered = await enforce();
const lastPruning = (
@ -234,28 +264,37 @@ it.each(["transaction", "iterator"] as const)(
)?.archivePruning;
expect(
recovered?.deferredReason,
JSON.stringify(
{
blockedCheckpoint: blocked.checkpoint,
recovered,
diagnosticsCount: diagnostics.length,
lastArchivePruning: {
completed: lastPruning?.completed,
checkpointCalls: lastPruning?.checkpointCalls,
checkpointIncomplete: lastPruning?.checkpointIncomplete,
checkpoint: lastPruning?.checkpoint,
walBytesBefore: lastPruning?.walBytesBefore,
walBytesAfter: lastPruning?.walBytesAfter,
},
},
// Checkpoint errors can contain paths; retain only the recorded health facts.
(key, value) => (key === "error" ? undefined : value),
),
recovered?.deferredReason === undefined
? undefined
: JSON.stringify(
{
blockedCheckpoint: blocked.checkpoint,
hostCheckpoint: database.walMaintenance.health,
now: Date.now(),
realClock: performance.timeOrigin + performance.now(),
dateNowMocked: vi.isMockFunction(Date.now),
recovered,
diagnosticsCount: diagnostics.length,
lastArchivePruning: {
completed: lastPruning?.completed,
checkpointCalls: lastPruning?.checkpointCalls,
checkpointIncomplete: lastPruning?.checkpointIncomplete,
checkpoint: lastPruning?.checkpoint,
walBytesBefore: lastPruning?.walBytesBefore,
walBytesAfter: lastPruning?.walBytesAfter,
},
},
// Checkpoint errors can contain paths; retain only the recorded health facts.
(key, value) => (key === "error" ? undefined : value),
),
).toBeUndefined();
expect(recovered?.totalBytesAfter).toBeLessThanOrEqual(maintenance.highWaterBytes!);
expect(diagnostics.at(-1)).toMatchObject({
archivePruning: { completed: true, checkpointIncomplete: 0 },
});
expect(() =>
JSON.stringify({ blocked, recovered, diagnostics, health: database.walMaintenance.health }),
).not.toThrow();
} finally {
writes.unsubscribe(observe);
release();

View file

@ -15,7 +15,8 @@ import {
import { runSqliteDeferredTransactionSync } from "./sqlite-transaction.js";
import {
onSqliteWalCheckpoint,
publishSqliteWalCheckpointHealth,
publishSqliteWalCheckpointObservation,
type SqliteWalCheckpointSnapshot,
} from "./sqlite-wal-checkpoint.js";
import { configureSqliteWalMaintenance } from "./sqlite-wal.js";
@ -285,10 +286,12 @@ describe("SQLite WAL checkpoint observations", () => {
busyTimeoutMs: 0,
});
const states: string[] = [];
const observations: SqliteWalCheckpointSnapshot[] = [];
const unsubscribe = onSqliteWalCheckpoint((observation) => {
if (observation.databasePath === databasePath) {
expect(maintenance.health?.state).toBe(observation.health.state);
states.push(observation.health.state);
observations.push({ health: observation.health, observedAtNs: observation.observedAtNs });
}
});
try {
@ -310,9 +313,9 @@ describe("SQLite WAL checkpoint observations", () => {
expect(() =>
assertNoActiveSqliteReaders(reader, "native transaction probe"),
).not.toThrow();
const observed = publishSqliteWalCheckpointHealth(databasePath, maintenance.health!);
expect(observed.activeReaders).toHaveLength(0);
expect(observed.readerDiagnostics).toHaveLength(1);
const observed = publishSqliteWalCheckpointObservation(databasePath, observations[0]!);
expect(observed.health.activeReaders).toHaveLength(0);
expect(observed.health.readerDiagnostics).toHaveLength(1);
},
{ operationLabel: "fixture.named-transaction" },
);

View file

@ -34,7 +34,8 @@ export type SqliteWalHealth = {
readerDiagnostics?: Array<Omit<SqliteReaderDiagnostics, "activeReaders">>;
};
export type SqliteWalCheckpointObservation = { databasePath: string; health: SqliteWalHealth };
export type SqliteWalCheckpointSnapshot = { health: SqliteWalHealth; observedAtNs: bigint };
export type SqliteWalCheckpointObservation = SqliteWalCheckpointSnapshot & { databasePath: string };
const checkpointListeners = resolveGlobalSingleton(
Symbol.for("openclaw.sqliteWalCheckpointListeners"),
() => new Set<(observation: SqliteWalCheckpointObservation) => void>(),
@ -79,12 +80,13 @@ function observeSqliteWalCheckpointHealth(
};
}
function notifyCheckpoint(databasePath: string, health: SqliteWalHealth): void {
function notifyCheckpoint(databasePath: string, snapshot: SqliteWalCheckpointSnapshot): void {
for (const listener of checkpointListeners) {
try {
listener({
databasePath: sqliteReaderDatabasePathKey(databasePath),
health: structuredClone(health),
health: structuredClone(snapshot.health),
observedAtNs: snapshot.observedAtNs,
});
} catch {
// Diagnostic consumers cannot change the native checkpoint's outcome.
@ -93,11 +95,14 @@ function notifyCheckpoint(databasePath: string, health: SqliteWalHealth): void {
}
/** Worker result transport relays the recorded fact and returns its enriched diagnostic snapshot. */
export function publishSqliteWalCheckpointHealth(
export function publishSqliteWalCheckpointObservation(
databasePath: string,
health: SqliteWalHealth,
): SqliteWalHealth {
const observed = observeSqliteWalCheckpointHealth(databasePath, health);
snapshot: SqliteWalCheckpointSnapshot,
): SqliteWalCheckpointSnapshot {
const observed = {
health: observeSqliteWalCheckpointHealth(databasePath, snapshot.health),
observedAtNs: snapshot.observedAtNs,
};
notifyCheckpoint(databasePath, observed);
return observed;
}
@ -128,7 +133,7 @@ export function createSqliteWalCheckpoint(
options: SqliteWalCheckpointOptions,
journalSizeLimitBytes: number,
) {
let health: SqliteWalHealth | undefined;
let snapshot: SqliteWalCheckpointSnapshot | undefined;
const checkpointObservation = (): SqliteWalHealth => ({
state: "error",
@ -137,7 +142,7 @@ export function createSqliteWalCheckpoint(
databaseBytes: null,
logFrames: null,
checkpointedFrames: null,
lastCompletedAtMs: health?.lastCompletedAtMs ?? null,
lastCompletedAtMs: snapshot?.health.lastCompletedAtMs ?? null,
consecutiveBlocked: 0,
warning: true,
});
@ -151,11 +156,14 @@ export function createSqliteWalCheckpoint(
warning: true,
error: formatErrorMessage(error),
};
health = options.databasePath
? observeSqliteWalCheckpointHealth(options.databasePath, failed)
: failed;
snapshot = {
observedAtNs: process.hrtime.bigint(),
health: options.databasePath
? observeSqliteWalCheckpointHealth(options.databasePath, failed)
: failed,
};
if (options.databasePath) {
notifyCheckpoint(options.databasePath, health);
notifyCheckpoint(options.databasePath, snapshot);
}
options.onCheckpointError?.(error);
};
@ -164,6 +172,8 @@ export function createSqliteWalCheckpoint(
mode: SqliteWalCheckpointMode,
row: Record<string, SQLOutputValue> | undefined,
): boolean => {
// Worker relays keep this same-process ordering fact even if the wall clock steps backward.
const observedAtNs = process.hrtime.bigint();
const observation = checkpointObservation();
let busy: boolean;
let sizeError: unknown;
@ -177,7 +187,7 @@ export function createSqliteWalCheckpoint(
if (observation.state === "complete") {
observation.lastCompletedAtMs = observation.observedAtMs;
} else {
observation.consecutiveBlocked = (health?.consecutiveBlocked ?? 0) + 1;
observation.consecutiveBlocked = (snapshot?.health.consecutiveBlocked ?? 0) + 1;
}
if (options.databasePath) {
try {
@ -196,11 +206,14 @@ export function createSqliteWalCheckpoint(
(observation.walBytes !== null &&
observation.databaseBytes !== null &&
observation.walBytes > Math.max(2 * observation.databaseBytes, journalSizeLimitBytes)));
health = options.databasePath
? observeSqliteWalCheckpointHealth(options.databasePath, observation)
: observation;
snapshot = {
observedAtNs,
health: options.databasePath
? observeSqliteWalCheckpointHealth(options.databasePath, observation)
: observation,
};
if (options.databasePath) {
notifyCheckpoint(options.databasePath, health);
notifyCheckpoint(options.databasePath, snapshot);
}
} catch (error) {
recordCheckpointError(error, observation);
@ -232,7 +245,10 @@ export function createSqliteWalCheckpoint(
);
},
get health() {
return health ? structuredClone(health) : undefined;
return snapshot ? structuredClone(snapshot.health) : undefined;
},
get snapshot() {
return snapshot ? structuredClone(snapshot) : undefined;
},
};
}

View file

@ -78,11 +78,18 @@ it("reports bounded page progress and retains the native checkpoint outcome", ()
maxPages: 31,
beforeMutation: () => observed.push(`before:${database.isTransaction}`),
onCommit: () => observed.push(`commit:${database.isTransaction}`),
afterCommit: () => observed.push(`settled:${database.isTransaction}`),
});
expect(observed).toEqual(["before:false", "before:true", "commit:true", "before:false"]);
expect(observed).toEqual([
"before:false",
"before:true",
"commit:true",
"settled:false",
"before:false",
]);
expect(result).toMatchObject({
checkpointCompleted: true,
checkpoint: { state: "complete", walBytes: 0 },
checkpoint: { health: { state: "complete", walBytes: 0 } },
checkpointCalls: 2,
checkpointIncomplete: 0,
vacuumPasses: 1,

View file

@ -3,18 +3,22 @@ import type { DatabaseSync } from "node:sqlite";
import { runWithSqliteBusyTimeout } from "./sqlite-busy-timeout.js";
import { isSqliteLockError } from "./sqlite-error-diagnostics.js";
import { runSqliteImmediateTransactionSync } from "./sqlite-transaction.js";
import type { SqliteWalCheckpointMode, SqliteWalHealth } from "./sqlite-wal-checkpoint.js";
import type {
SqliteWalCheckpointMode,
SqliteWalCheckpointSnapshot,
} from "./sqlite-wal-checkpoint.js";
export type SqliteWalReclamationOptions = {
maxPages?: number;
checkpointMode?: SqliteWalCheckpointMode;
beforeMutation?: () => void;
onCommit?: () => void;
afterCommit?: () => void;
};
export type SqliteWalReclamationResult = {
checkpointCompleted: boolean;
checkpoint?: SqliteWalHealth;
checkpoint?: SqliteWalCheckpointSnapshot;
freePagesBefore: number | null;
remainingFreePages: number | null;
checkpointCalls: number;
@ -114,6 +118,7 @@ export function reclaimSqliteWalFreePages(
} finally {
result.vacuumMs += performance.now() - startedAt;
}
options.afterCommit?.();
if (checkpoint()) {
result.remainingFreePages = freePages();
}

View file

@ -557,7 +557,7 @@ export function configureSqliteWalMaintenance(
if (!result) {
throw new Error("SQLite page reclamation owner is unavailable");
}
return { ...result, checkpoint: checkpointOwner.health };
return { ...result, checkpoint: checkpointOwner.snapshot };
};
let timer: IntervalHandle | null = null;