mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 17:53:39 +00:00
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:
parent
aca0c60535
commit
bef4401fbc
20 changed files with 444 additions and 93 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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">;
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
);
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 };
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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 };
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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" },
|
||||
);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue