fix: keep session writes responsive during archive cleanup (#156014)

* fix: keep session writes responsive during archive cleanup

* fix: load archive history reader lazily

* fix: resume budget cleanup after pooled worker checkpoints

* fix: keep archive maintenance fixtures on their host owner
This commit is contained in:
Peter Steinberger 2026-09-22 23:09:35 -07:00 • committed by GitHub
parent cb6ef2b3e9
commit b981f930dc
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
40 changed files with 2223 additions and 1166 deletions

View file

@ -1668,17 +1668,41 @@ 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
section. Each unit checks current authority before checkpointing and vacuum,
authorizes commit, and joins native settlement. Cache eviction between units can
refresh the host claim only for the same physical database; no dispatched mutation
is replayed. Incognito maintenance retains its in-process owner.
Archive-row and unpublished-name reads follow
validation. After removing a derived archive file, pruning reacquires before the
canonical row-deletion transaction; an acquisition failure propagates without
deleting that recovery row.
Archive pruning retains its maintenance and archive-worker lifetimes while each
page-reclamation unit acquires and releases the physical writer separately.
Foreground session writes can run between units. Durable archive metadata reads
use the existing history worker; conditional deletions, legacy file removal, and
page reclamation use the agent database execution broker.
The host captures the physical file before
waiting and rechecks the original path alias and live authority before effects
and at worker admission. Each write joins native settlement without replaying a
dispatched mutation. Host connection eviction does not redirect work or require
a synchronous database reopen. Process-held incognito maintenance retains its
existing in-process owner and remains a separate worker migration.
Canonical archive removal holds one writer section through selection, derived-file
removal, and conditional row deletion. It rechecks pressure, authority, and the
selected published row before deleting the canonical recovery copy. A failed
admission or changed row preserves that copy and stops the current
cleanup attempt. Legacy file removal checks exact canonical filename ownership,
stats the file, and unlinks it in one synchronous worker write transaction. It
preserves files owned by published or unpublished rows and holds the SQLite writer
lock through unlink. Since file removal cannot roll back, the host grants commit
immediately before unlink; no-effect outcomes also require current commit authority.
Native settlement finishes before the item writer is released. Filesystem inventory
and successive page drains run outside the item writer.
Aggregate pruning diagnostics report the whole operation separately from actual
writer waits.
Archive order, retention policy, schemas, and update behavior are unchanged.
Worker retirement preserves the original operation failure without reporting it
again as a cleanup failure. A successfully retired execution owner is released for
later requests; genuine native-close and lease-cleanup failures retain their
existing retry custody.
Successful pooled-agent close relays its recorded WAL checkpoint after native and
lease cleanup settle. The original generation and physical database identities
fence that observation, and the budget owner releases deferral only for a newer
completed checkpoint.
Usage-cache rollup writes, pruning, and refresh-lock changes use the same async
agent-database admission. A cold mutation waits for the existing integrity worker;

View file

@ -165,18 +165,19 @@ it.each([
},
});
}
// The failed run already attempted canonical cleanup. Client close only
// releases its borrow; the canonical resource retains the failed retirement.
const failed = worker;
// The client can borrow a fresh executor after native retirement; a stale
// hash proves admission without replaying the uncertain delete.
await expect(
Promise.resolve().then(() =>
failed.run(
async () => undefined,
() => undefined,
),
worker.run(
(scope) =>
scope.execute({
...command,
input: { ...command.input, expectedHash: "not-current" },
}),
() => undefined,
),
).rejects.toThrow("admission is closed");
await failed.close();
).resolves.toEqual({ ok: true, value: false });
await worker.close();
worker = undefined;
await closeOpenClawAgentDatabasesAsync(state.stateDir);
const recovered = openOpenClawAgentDatabase(options);

View file

@ -8,7 +8,6 @@ import {
} from "../../../config/sessions/session-accessor.js";
import {
resolveSqliteReadScope,
runExclusiveSqliteSessionWrite,
toDatabaseOptions,
} from "../../../config/sessions/session-accessor.sqlite-scope.js";
import { pruneAllSessionTranscriptArchivesToHighWater } from "../../../config/sessions/session-history-archive-pruning.js";
@ -128,18 +127,12 @@ it("submits deferred child results after canonical archive pruning without poiso
});
const scope = resolveSqliteReadScope(archived.target);
const pruned = await runExclusiveSqliteSessionWrite(
scope,
() =>
pruneAllSessionTranscriptArchivesToHighWater({
archiveDirectory: path.dirname(archivePath),
databaseOptions: toDatabaseOptions(scope),
highWaterBytes: 0,
storePath: archived.target.storePath,
withArchiveWrite: (write) => write(),
}),
"session.history.archive-prune",
);
const pruned = await pruneAllSessionTranscriptArchivesToHighWater({
archiveDirectory: path.dirname(archivePath),
databaseOptions: toDatabaseOptions(scope),
highWaterBytes: 0,
storePath: archived.target.storePath,
});
expect(pruned.removedFiles).toBe(1);
await expect(access(archivePath)).rejects.toMatchObject({ code: "ENOENT" });
expect(

View file

@ -1,5 +1,6 @@
import fs from "node:fs";
import path from "node:path";
import { err, ok, type Result } from "@openclaw/normalization-core/result";
import { resolveRealpathOrAbsolute as canonicalizePathForComparison } from "../../infra/boundary-path.js";
import { runTasksWithConcurrency } from "../../utils/run-with-concurrency.js";
import { isMigrationArchiveArtifactName } from "./artifacts.js";
@ -23,6 +24,52 @@ export type SessionsDirFileStat = {
const SESSIONS_DIR_STAT_CONCURRENCY = 8;
// A removed empty file is success; bytes alone cannot signal removal.
export type FileRemovalResult = Result<number, "not-removed">;
export async function removeFileIfExists(filePath: string): Promise<FileRemovalResult> {
const stat = await fs.promises.stat(filePath).catch(() => null);
if (!stat?.isFile()) {
return err("not-removed");
}
// Forced removal would count paths another cleanup already removed after stat.
return fs.promises.rm(filePath).then(
() => ok(stat.size),
() => err("not-removed"),
);
}
export async function removeFileForBudget(params: {
filePath: string;
canonicalPath?: string;
dryRun: boolean;
fileSizesByPath: Map<string, number>;
simulatedRemovedPaths: Set<string>;
onRemovedPath?: (canonicalPath: string) => void;
}): Promise<FileRemovalResult> {
const resolvedPath = path.resolve(params.filePath);
const canonicalPath = params.canonicalPath ?? canonicalizePathForComparison(resolvedPath);
if (params.dryRun) {
// Dry-run deletion is path-deduped so a transcript and pointer alias cannot count the same
// artifact twice against the simulated budget.
if (params.simulatedRemovedPaths.has(canonicalPath)) {
return err("not-removed");
}
const size = params.fileSizesByPath.get(canonicalPath);
if (size === undefined) {
return err("not-removed");
}
params.simulatedRemovedPaths.add(canonicalPath);
params.onRemovedPath?.(canonicalPath);
return ok(size);
}
const removal = await removeFileIfExists(resolvedPath);
if (removal.ok) {
params.onRemovedPath?.(canonicalPath);
}
return removal;
}
export async function readSessionsDirFiles(sessionsDir: string): Promise<SessionsDirFileStat[]> {
const dirEntries = await fs.promises
.readdir(sessionsDir, { withFileTypes: true })

View file

@ -1,7 +1,7 @@
// Session disk-budget enforcement prunes orphaned artifacts before deleting store entries.
import fs from "node:fs";
import path from "node:path";
import { err, ok, type Result } from "@openclaw/normalization-core/result";
import { err } from "@openclaw/normalization-core/result";
import {
normalizeLowercaseStringOrEmpty,
normalizeOptionalLowercaseString,
@ -24,6 +24,9 @@ import {
isSessionPromptBlobTempArtifactName,
readSessionPromptBlobFiles,
readSessionsDirFiles,
removeFileForBudget,
removeFileIfExists,
type FileRemovalResult,
type SessionPhysicalDiskUsage,
type SessionsDirFileStat,
} from "./disk-budget-files.js";
@ -36,6 +39,7 @@ import { readLegacyCompactionSnapshotPaths } from "./legacy-compaction-history.j
import { resolveSessionArtifactDirectory, resolveSessionFilePathCore } from "./paths.js";
import type { SqliteSessionArchivePruningDiagnostics } from "./session-accessor.sqlite-contract.js";
import { timeArchivePruningAsync } from "./session-history-archive-pruning-diagnostics.js";
import type { SessionLegacyArchiveRemovalResult } from "./session-history-archive-pruning.types.js";
import { projectSessionStoreForPersistence } from "./skill-prompt-blobs.js";
import { isSessionEntryDiskBudgetEvictable } from "./store-maintenance.js";
import type { SessionEntry } from "./types.js";
@ -202,9 +206,9 @@ export async function hasRetainedSessionTranscriptArchives(storePath: string): P
/** Removes oldest retained archives and legacy compact backups, remeasuring after each file. */
export async function pruneSessionTranscriptArchivesToHighWater(params: {
diagnostics?: SqliteSessionArchivePruningDiagnostics;
excludeNames?: ReadonlySet<string>;
highWaterBytes: number;
storePath: string;
removeFile?: (file: SessionsDirFileStat) => Promise<SessionLegacyArchiveRemovalResult>;
}): Promise<{ removedFiles: number; usage: SessionPhysicalDiskUsage }> {
// Oldest-first is the hard-cap sacrifice order: under extreme pressure this
// may prune an archive the current pass just extracted, which is preferred
@ -212,10 +216,7 @@ export async function pruneSessionTranscriptArchivesToHighWater(params: {
const { diagnostics } = params;
const files = await timeArchivePruningAsync(diagnostics, "legacyInventoryMs", async () =>
(await readSessionsDirFiles(resolveSessionArtifactDirectory(params.storePath)))
.filter(
(file) =>
isRetainedSessionTranscriptArchiveName(file.name) && !params.excludeNames?.has(file.name),
)
.filter((file) => isRetainedSessionTranscriptArchiveName(file.name))
.toSorted((left, right) => left.mtimeMs - right.mtimeMs),
);
let usage = await timeArchivePruningAsync(diagnostics, "measurementMs", () =>
@ -226,21 +227,26 @@ export async function pruneSessionTranscriptArchivesToHighWater(params: {
if (usage.totalBytes <= params.highWaterBytes) {
break;
}
if (
!(
await timeArchivePruningAsync(diagnostics, "fileRemovalMs", () =>
removeFileIfExists(file.path),
)
).ok
) {
const removal = params.removeFile
? await params.removeFile(file)
: (
await timeArchivePruningAsync(diagnostics, "fileRemovalMs", () =>
removeFileIfExists(file.path),
)
).ok
? "removed"
: "failed";
if (removal === "failed") {
if (diagnostics) {
diagnostics.failedRemovals = (diagnostics.failedRemovals ?? 0) + 1;
}
continue;
}
removedFiles += 1;
if (diagnostics) {
diagnostics.removedFiles = (diagnostics.removedFiles ?? 0) + 1;
if (removal === "removed") {
removedFiles += 1;
if (diagnostics) {
diagnostics.removedFiles = (diagnostics.removedFiles ?? 0) + 1;
}
}
usage = await timeArchivePruningAsync(diagnostics, "measurementMs", () =>
measureSessionPhysicalDiskUsage(params.storePath),
@ -318,52 +324,6 @@ function isDiskBudgetRemovableSessionFile(
);
}
// A removed empty file is success; bytes alone cannot signal removal.
type FileRemovalResult = Result<number, "not-removed">;
async function removeFileIfExists(filePath: string): Promise<FileRemovalResult> {
const stat = await fs.promises.stat(filePath).catch(() => null);
if (!stat?.isFile()) {
return err("not-removed");
}
// Forced removal would count paths another cleanup already removed after stat.
return fs.promises.rm(filePath).then(
() => ok(stat.size),
() => err("not-removed"),
);
}
async function removeFileForBudget(params: {
filePath: string;
canonicalPath?: string;
dryRun: boolean;
fileSizesByPath: Map<string, number>;
simulatedRemovedPaths: Set<string>;
onRemovedPath?: (canonicalPath: string) => void;
}): Promise<FileRemovalResult> {
const resolvedPath = path.resolve(params.filePath);
const canonicalPath = params.canonicalPath ?? canonicalizePathForComparison(resolvedPath);
if (params.dryRun) {
// Dry-run deletion is path-deduped so a transcript and pointer alias cannot count the same
// artifact twice against the simulated budget.
if (params.simulatedRemovedPaths.has(canonicalPath)) {
return err("not-removed");
}
const size = params.fileSizesByPath.get(canonicalPath);
if (size === undefined) {
return err("not-removed");
}
params.simulatedRemovedPaths.add(canonicalPath);
params.onRemovedPath?.(canonicalPath);
return ok(size);
}
const removal = await removeFileIfExists(resolvedPath);
if (removal.ok) {
params.onRemovedPath?.(canonicalPath);
}
return removal;
}
async function removePromptBlobFileForBudget(params: {
file: SessionsDirFileStat;
projectedPromptBlobRefCounts: ReadonlyMap<string, number>;

View file

@ -3,6 +3,7 @@ import fs from "node:fs/promises";
import path from "node:path";
import { describe, expect, it, vi } from "vitest";
import { withTestDir } from "../../test-helpers/temp-dir.js";
import { removeFileIfExists } from "./disk-budget-files.js";
import {
enforceSessionDiskBudget,
measureSessionPhysicalDiskUsage,
@ -148,7 +149,16 @@ it("counts empty retained archives under pressure and returns real disk usage",
const excludedName = `keep.jsonl.deleted.${ARCHIVE_STAMP}`;
const excluded = await writeOldFile(dir, excludedName);
await fs.writeFile(path.join(dir, "filler.bin"), Buffer.alloc(128));
const params = { storePath, highWaterBytes: 64, excludeNames: new Set([excludedName]) };
const params: Parameters<typeof pruneSessionTranscriptArchivesToHighWater>[0] = {
storePath,
highWaterBytes: 64,
removeFile: async (file) => {
if (file.name === excludedName) {
return "preserved";
}
return (await removeFileIfExists(file.path)).ok ? "removed" : "failed";
},
};
const result = await pruneSessionTranscriptArchivesToHighWater(params);

View file

@ -78,9 +78,6 @@ export type SqliteSessionArtifactPreparationDiagnostics =
/** One pruning attempt retains only aggregate stage observations. */
export type SqliteSessionArchivePruningDiagnostics = {
trigger: "initial" | "after-eviction" | "final";
admissionMs?: number;
cachedAdmissions?: number;
asyncAdmissions?: number;
checkpointCalls?: number;
checkpointIncomplete?: number;
checkpoint?: SqliteWalHealth;
@ -107,7 +104,6 @@ export type SqliteSessionArchivePruningDiagnostics = {
export type SqliteSessionWriteDiagnostics = SqliteSessionReclamationDiagnostics & {
artifactPreparation?: SqliteSessionArtifactPreparationDiagnostics;
archivePruning?: SqliteSessionArchivePruningDiagnostics;
reclamationAdmission?: SqliteSessionReclamationAdmissionDiagnostics;
};

View file

@ -0,0 +1,245 @@
import { createHash } from "node:crypto";
import fs from "node:fs";
import path from "node:path";
import { afterEach, expect, it, vi } from "vitest";
import { createDeferred, withTestTimeout } from "../../../test/helpers/promise.js";
import { executeSqliteQuerySync } from "../../infra/kysely-sync.js";
import type { SqliteWalReclamationResult } from "../../infra/sqlite-wal-reclamation.js";
import {
openOpenClawAgentDatabase,
type OpenClawAgentDatabase,
} from "../../state/openclaw-agent-db.js";
import { ensureSessionTranscriptArchiveSchema } from "../../state/openclaw-agent-session-transcript-archive-schema.js";
import * as writeAdmission from "../../state/openclaw-agent-write-admission.js";
import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js";
import { resolveRegisteredSqliteTranscriptArchiveName } from "./session-accessor.sqlite-archive-artifact.js";
import { withSqliteSessionPageReclamation } from "./session-accessor.sqlite-page-reclamation.js";
import * as pageReclamation from "./session-accessor.sqlite-page-reclamation.js";
import {
getSessionKysely,
runExclusiveSqliteSessionWrite,
} from "./session-accessor.sqlite-scope.js";
import { pruneAllSessionTranscriptArchivesToHighWater } from "./session-history-archive-pruning.js";
afterEach(() => vi.restoreAllMocks());
function seedArchive(database: OpenClawAgentDatabase, archiveDirectory: string) {
const sessionId = "binding-retained-history";
const generation = "binding-generation";
const createdAt = 1;
const archiveName = resolveRegisteredSqliteTranscriptArchiveName({
createdAt,
encoding: "identity",
generation,
reason: "deleted",
sessionId,
});
const content = Buffer.from("synthetic retained history\n");
ensureSessionTranscriptArchiveSchema(database.db);
executeSqliteQuerySync(
database.db,
getSessionKysely(database.db)
.insertInto("session_transcript_archives")
.values({
archive_blob: content,
archive_name: archiveName,
archive_sha256: createHash("sha256").update(content).digest("hex"),
created_at: createdAt,
encoding: "identity",
generation,
published_at: createdAt,
reason: "deleted",
session_id: sessionId,
session_key: "agent:main:binding-history",
}),
);
fs.mkdirSync(archiveDirectory, { recursive: true });
const archivePath = path.join(archiveDirectory, archiveName);
fs.writeFileSync(archivePath, content);
return { archivePath, content, sessionId };
}
function readArchives(database: OpenClawAgentDatabase) {
return executeSqliteQuerySync(
database.db,
getSessionKysely(database.db)
.selectFrom("session_transcript_archives")
.select(["archive_name", "archive_sha256", "generation", "published_at", "session_id"]),
).rows;
}
it.runIf(process.platform !== "win32").each([false, true])(
"keeps archive pruning on its captured database (retargeted alias: %s)",
async (retarget) => {
await withOpenClawTestState(
{ prefix: "page-reclamation-binding-", scenario: "minimal", layout: "state-only" },
async (state) => {
const options = { agentId: "main", env: state.env };
const original = openOpenClawAgentDatabase(options);
const replacement = retarget
? openOpenClawAgentDatabase({
...options,
path: path.join(state.root, "replacement.sqlite"),
})
: undefined;
const archiveDatabase = replacement ?? original;
const archiveDirectory = state.sessionsDir();
const archive = seedArchive(archiveDatabase, archiveDirectory);
const before = readArchives(archiveDatabase);
const alias = path.join(state.root, "alias.sqlite");
fs.symlinkSync(original.path, alias);
const aliasOptions = { ...options, path: alias };
const withPages = pageReclamation.withSqliteSessionPageReclamation;
let captured = false;
vi.spyOn(pageReclamation, "withSqliteSessionPageReclamation").mockImplementation(
<T>(...args: Parameters<typeof withPages<T>>) => {
const [input, run] = args;
return withPages(input, (...context) => {
captured = true;
if (replacement) {
fs.unlinkSync(alias);
fs.symlinkSync(replacement.path, alias);
}
return run(...context);
});
},
);
let failure: unknown;
let removedFiles: number | undefined;
try {
const result = await pruneAllSessionTranscriptArchivesToHighWater({
archiveDirectory,
databaseOptions: aliasOptions,
highWaterBytes: 0,
storePath: alias,
});
removedFiles = result.removedFiles;
} catch (error) {
failure = error;
}
expect(captured).toBe(true);
if (retarget) {
expect(readArchives(archiveDatabase)).toEqual(before);
expect(fs.readFileSync(archive.archivePath)).toEqual(archive.content);
expect(failure).toBeInstanceOf(Error);
expect(String(failure)).toContain("file identity changed");
} else {
expect(failure).toBeUndefined();
expect(removedFiles).toBe(1);
expect(readArchives(archiveDatabase)).toEqual([]);
expect(fs.existsSync(archive.archivePath)).toBe(false);
}
},
);
},
);
it.runIf(process.platform !== "win32").each([false, true])(
"revalidates the original alias after worker write admission queues (retargeted alias: %s)",
async (retarget) => {
let fixtureRoot = "";
await withOpenClawTestState(
{ prefix: "page-reclamation-queued-binding-", scenario: "minimal", layout: "state-only" },
async (state) => {
fixtureRoot = state.root;
const options = { agentId: "main", env: state.env };
const original = openOpenClawAgentDatabase(options);
const replacement = openOpenClawAgentDatabase({
...options,
path: path.join(state.root, "replacement.sqlite"),
});
const originalArchive = seedArchive(original, path.join(state.root, "original-archives"));
const replacementArchive = seedArchive(
replacement,
path.join(state.root, "replacement-archives"),
);
const originalRows = readArchives(original);
const replacementRows = readArchives(replacement);
const pageSize = Number(original.db.prepare("PRAGMA page_size").get()?.page_size);
// sqlite-allow-raw -- Synthetic free pages prove whether the real queued vacuum was authorized.
original.db
.prepare(
"INSERT INTO cache_entries(scope, key, blob, updated_at) VALUES (?, ?, zeroblob(?), 1)",
)
.run("queued-binding", "free-pages", pageSize * 1024);
original.db.prepare("DELETE FROM cache_entries WHERE scope = ?").run("queued-binding");
expect(original.walMaintenance.checkpoint()).toBe(true);
const readFreePages = () =>
Number(original.db.prepare("PRAGMA freelist_count").get()?.freelist_count);
const freePagesBefore = readFreePages();
expect(freePagesBefore).toBeGreaterThan(512);
const alias = path.join(state.root, "alias.sqlite");
fs.symlinkSync(original.path, alias);
const aliasOptions = { ...options, path: alias };
const requested = createDeferred();
let observing = false;
let preparedPath: string | undefined;
const write = writeAdmission.runOpenClawAgentWorkerWrite;
vi.spyOn(writeAdmission, "runOpenClawAgentWorkerWrite").mockImplementation(
<T>(...args: Parameters<typeof write<T>>) => {
const result = write(...args);
if (observing && args[0].path === preparedPath) {
requested.resolve();
}
return result;
},
);
let failure: unknown;
let result: SqliteWalReclamationResult | undefined;
try {
result = await withSqliteSessionPageReclamation(
aliasOptions,
async (reclaimPages, _assertCurrent, preparedOptions) => {
preparedPath = preparedOptions.path;
const entered = createDeferred();
const release = createDeferred();
const blocker = runExclusiveSqliteSessionWrite(
preparedOptions,
async () => {
entered.resolve();
await release.promise;
},
"session.transcript.batch",
);
void blocker.catch(entered.reject);
let reclamation: Promise<SqliteWalReclamationResult> | undefined;
try {
await entered.promise;
observing = true;
reclamation = reclaimPages(1);
void reclamation.catch(() => {});
await withTestTimeout(requested.promise, 10_000, "Page write was not queued");
if (retarget) {
fs.unlinkSync(alias);
fs.symlinkSync(replacement.path, alias);
}
release.resolve();
return await reclamation;
} finally {
observing = false;
release.resolve();
await Promise.allSettled([blocker, reclamation]);
}
},
);
} catch (error) {
failure = error;
}
expect(readArchives(original)).toEqual(originalRows);
expect(readArchives(replacement)).toEqual(replacementRows);
expect(fs.readFileSync(originalArchive.archivePath)).toEqual(originalArchive.content);
expect(fs.readFileSync(replacementArchive.archivePath)).toEqual(replacementArchive.content);
if (retarget) {
expect(readFreePages()).toBe(freePagesBefore);
expect(failure).toBeInstanceOf(Error);
expect(String(failure)).toContain("file identity changed");
} else {
expect(failure).toBeUndefined();
expect(result).toMatchObject({ checkpointCompleted: true, vacuumPasses: 1 });
expect(readFreePages()).toBeLessThan(freePagesBefore);
}
},
);
expect(fs.existsSync(fixtureRoot)).toBe(false);
},
);

View file

@ -1,112 +1,298 @@
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";
import {
assertExistingDatabaseIdentity,
readDatabasePathIdentitySync,
} from "../../infra/sqlite-worker-identity.js";
import { createSqliteWorkerOperationAdmission } from "../../infra/sqlite-worker-operation-admission.js";
import type { SqliteWorkerStore } from "../../infra/sqlite-worker-store.js";
import {
isIncognitoOpenClawAgentSqlitePath,
type OpenClawAgentDatabaseOptions,
} from "../../state/openclaw-agent-db.js";
import { runPreparedSqliteSessionReclamation } from "./session-accessor.sqlite-reclamation-run.js";
import { withSqliteReclamationWorker } from "./session-accessor.sqlite-reclamation-worker.js";
import type {
AgentDatabaseExecutionFileIdentity,
AgentDatabaseOperations,
AgentDatabaseRequestExecutionSource,
} from "../../state/openclaw-agent-execution-contract.js";
import {
captureOpenClawAgentDatabaseExecution,
supportsOpenClawAgentDatabaseExecution,
} from "../../state/openclaw-agent-execution.js";
import { runExclusiveSqliteTranscriptArchiveWorker } from "./session-accessor.sqlite-archive.js";
import type { ReclamationDatabaseOptions } from "./session-accessor.sqlite-lifecycle-types.js";
import { resolveSessionReclamationDatabaseOptions } from "./session-accessor.sqlite-reclamation.js";
import {
runExclusiveSqliteSessionWrite,
withSqliteSessionDatabase,
} from "./session-accessor.sqlite-scope.js";
import { withSqliteMutationWorkerLifetime } from "./session-accessor.sqlite-worker-request.js";
import type {
PublishedSessionTranscriptArchive,
SessionArchivePruningOperations,
} from "./session-history-archive-pruning.types.js";
export type SqliteSessionPageReclaimer = (maxPages?: number) => Promise<SqliteWalReclamationResult>;
/** Acquire the archive worker before the caller's writer, so retained work cannot deadlock it. */
export async function withSqliteSessionPageReclamation<T>(
/** Preview reads share the history reader without waiting for archive or writer admission. */
export async function readSqliteSessionArchivePruning(
input: OpenClawAgentDatabaseOptions,
run: (reclaim: SqliteSessionPageReclaimer) => Promise<T>,
): Promise<T> {
): Promise<PublishedSessionTranscriptArchive | null> {
const options = resolveSessionReclamationDatabaseOptions(input);
if (isIncognitoOpenClawAgentSqlitePath(options.path, options)) {
return run(async (maxPages) =>
withSqliteSessionDatabase(options, (database) =>
database.walMaintenance.reclaimFreePages({ maxPages }),
),
if (!supportsOpenClawAgentDatabaseExecution(options)) {
const { readSessionArchivePruningInDatabase } =
await import("./session-history-archive-pruning.worker.js");
return withSqliteSessionDatabase(options, (database) =>
readSessionArchivePruningInDatabase(database),
);
}
const physical = readDatabasePathIdentitySync(options.path);
if (!physical.key.startsWith("file:")) {
return null;
}
const databaseOptions = { ...options, path: physical.canonicalPath };
const expectedIdentity: AgentDatabaseExecutionFileIdentity = {
kind: "file",
physicalIdentity: physical.key.slice("file:".length),
nativeLocation: physical.canonicalPath,
};
const { withSessionHistoryWorkerDatabase } =
await import("./session-transcript-worker-runtime.js");
assertExistingDatabaseIdentity(options.path, physical.key);
return withSessionHistoryWorkerDatabase(databaseOptions, async (reader) => {
assertExistingDatabaseIdentity(options.path, physical.key);
const result = await reader.readArchivePruning({
env: databaseOptions.env,
expectedIdentity,
});
reader.assertCurrent();
assertExistingDatabaseIdentity(options.path, physical.key);
return result;
});
}
/** Retain source custody before archive admission; each page unit owns its writer separately. */
export async function withSqliteSessionPageReclamation<T>(
input: OpenClawAgentDatabaseOptions,
run: (
reclaim: SqliteSessionPageReclaimer,
assertCurrent: () => void,
databaseOptions: ReclamationDatabaseOptions,
archives: SessionArchivePruningOperations,
) => Promise<T>,
): Promise<T> {
const options = resolveSessionReclamationDatabaseOptions(input);
const incognito = isIncognitoOpenClawAgentSqlitePath(options.path, options);
const nativeOwner = !supportsOpenClawAgentDatabaseExecution(options);
// Capture the original alias before admission can wait. Workers use only this physical path.
const physical = incognito ? undefined : readDatabasePathIdentitySync(options.path);
if (physical && !physical.key.startsWith("file:")) {
throw new Error("SQLite archive pruning requires its existing database");
}
return withSqliteMutationWorkerLifetime(options, async ({ assertCurrent, signal }) => {
const retained = await runExclusiveSqliteSessionWrite(
options,
async () => {
if (nativeOwner) {
// Incognito and explicit Doctor/cleanup maintenance retain their existing native owner.
const {
readSessionArchivePruningInDatabase,
deletePublishedSessionArchiveInDatabase,
removeLegacySessionArchiveInDatabase,
} = await import("./session-history-archive-pruning.worker.js");
const databaseOptions = physical ? { ...options, path: physical.canonicalPath } : options;
const assertNativeCurrent = () => {
assertCurrent();
return retainOpenClawAgentDatabaseReadOnly(options);
},
"session.reclamation.retain",
);
if (!retained.found) {
throw new Error("SQLite page reclamation lost its prepared database");
if (physical) {
assertExistingDatabaseIdentity(options.path, physical.key);
assertExistingDatabaseIdentity(databaseOptions.path, physical.key);
}
};
const runNative = () => {
assertNativeCurrent();
return run(
(maxPages) =>
runExclusiveSqliteSessionWrite(
databaseOptions,
async () =>
withSqliteSessionDatabase(
databaseOptions,
(database) => {
assertNativeCurrent();
return database.walMaintenance.reclaimFreePages({
maxPages,
beforeMutation: assertNativeCurrent,
onCommit: assertNativeCurrent,
});
},
assertNativeCurrent,
),
"session.history.free-pages",
),
assertNativeCurrent,
databaseOptions,
{
read: async () =>
await withSqliteSessionDatabase(
databaseOptions,
(database) => {
assertNativeCurrent();
return readSessionArchivePruningInDatabase(database);
},
assertNativeCurrent,
),
removeLegacy: async (filePath) =>
await withSqliteSessionDatabase(
databaseOptions,
(database) =>
removeLegacySessionArchiveInDatabase(
database,
databaseOptions,
filePath,
assertNativeCurrent,
),
assertNativeCurrent,
),
deletePublished: async (row) =>
await withSqliteSessionDatabase(
databaseOptions,
(database) => {
deletePublishedSessionArchiveInDatabase(
database,
databaseOptions,
row,
assertNativeCurrent,
);
},
assertNativeCurrent,
),
},
);
};
return physical ? runExclusiveSqliteTranscriptArchiveWorker(runNative, signal) : runNative();
}
let { database, claim } = retained;
const physicalIdentity = claim.identity;
const databaseOptions = {
...options,
path: readOpenClawAgentDatabaseIdentity(database).filename,
if (!physical) {
throw new Error("SQLite archive pruning requires its existing file owner");
}
const databaseOptions = { ...options, path: physical.canonicalPath };
const expectedIdentity: AgentDatabaseExecutionFileIdentity = {
kind: "file",
physicalIdentity: physical.key.slice("file:".length),
nativeLocation: physical.canonicalPath,
};
const execution = captureOpenClawAgentDatabaseExecution(databaseOptions, { expectedIdentity });
const assertPruningCurrent = () => {
assertCurrent();
execution.assertCurrent();
assertExistingDatabaseIdentity(options.path, physical.key);
assertExistingDatabaseIdentity(databaseOptions.path, physical.key);
};
const source: AgentDatabaseRequestExecutionSource = {
assertCurrent: assertPruningCurrent,
createAdmission(binding) {
return () => ({
nativeLocations: binding.nativeLocations,
admission: createSqliteWorkerOperationAdmission((request, grant) => {
binding.authorize(request);
assertPruningCurrent();
if (!grant()) {
throw new Error("SQLite archive pruning authority expired");
}
}),
});
},
};
const write = async <Value>(
operation: (
worker: Pick<SqliteWorkerStore<AgentDatabaseOperations>, "execute">,
) => Promise<Value>,
label: "session.history.free-pages" | "session.history.archive-prune",
): Promise<Value> => {
const result = await runExclusiveSqliteSessionWrite(
databaseOptions,
async () => {
assertPruningCurrent();
return execution.runExisting(
source,
async (worker) => ({ value: await operation(worker) }),
{
retireNativeOnFailure: true,
},
);
},
label,
undefined,
"worker",
);
assertPruningCurrent();
if (!result) {
throw new Error("SQLite archive pruning lost its prepared database");
}
return result.value;
};
try {
return await withSqliteReclamationWorker(
const [{ withSessionHistoryWorkerDatabase }, { maintenanceLane }] = await Promise.all([
import("./session-transcript-worker-runtime.js"),
import("./session-transcript-worker-resources.js"),
]);
assertPruningCurrent();
return await withSessionHistoryWorkerDatabase(
databaseOptions,
claim,
(worker) =>
run((maxPages) =>
withSqliteMutationWorkerLifetime(databaseOptions, async (request) => {
const assertRequestCurrent = () => {
assertCurrent();
request.assertCurrent();
};
assertRequestCurrent();
if (!claim.isCurrent()) {
// Archive-file I/O may retire the borrowed host handle; keep the same physical store.
const reopened = retainOpenClawAgentDatabaseReadOnly(databaseOptions);
if (!reopened.found) {
throw new Error("SQLite page reclamation lost its prepared database");
}
if (reopened.claim.identity !== physicalIdentity) {
reopened.claim.release();
throw new Error("SQLite page reclamation database path was replaced");
}
claim.release();
({ database, claim } = reopened);
}
const result = await runPreparedSqliteSessionReclamation(
{
plan: {
kind: "maintenance-pages",
databaseOptions,
materializedPlans: [],
maxPages,
},
},
{
database,
claim,
worker,
commitGate: request.commitGate,
assertRequestCurrent,
},
);
if (result.kind !== "maintenance-pages") {
throw new Error("SQLite page reclamation returned another operation's result");
}
if (result.value.checkpoint) {
result.value.checkpoint = publishSqliteWalCheckpointObservation(
databaseOptions.path,
result.value.checkpoint,
async (reader) =>
await runExclusiveSqliteTranscriptArchiveWorker(async () => {
assertPruningCurrent();
return run(
async (maxPages) => {
const result = await write(
(worker) =>
worker.execute({
type: "session.archivePruning.reclaimPages",
input: { maxPages },
}),
"session.history.free-pages",
);
}
return result.value;
}),
),
assertCurrent,
signal,
if (result.checkpoint) {
result.checkpoint = publishSqliteWalCheckpointObservation(
databaseOptions.path,
result.checkpoint,
);
}
return result;
},
assertPruningCurrent,
databaseOptions,
{
read: async () => {
assertPruningCurrent();
const result = await reader.readArchivePruning({
env: databaseOptions.env,
expectedIdentity,
});
assertPruningCurrent();
return result;
},
removeLegacy: (filePath) =>
write(
(worker) =>
worker.execute({
type: "session.archivePruning.removeLegacy",
input: { filePath },
}),
"session.history.archive-prune",
),
deletePublished: (row) =>
write(
(worker) =>
worker.execute({
type: "session.archivePruning.deletePublished",
input: row,
}),
"session.history.archive-prune",
),
},
);
}, signal),
maintenanceLane,
);
} finally {
claim.release();
await execution.release();
}
});
}

View file

@ -5,22 +5,28 @@ 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 * 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 * as workerAdmission from "../../infra/sqlite-worker-operation-admission.js";
import { createDeferredCore } from "../../shared/deferred.js";
import { closeCachedOpenClawAgentDatabase } from "../../state/openclaw-agent-db-lifecycle.js";
import { withOpenClawAgentDatabaseReadOnly } from "../../state/openclaw-agent-db-readonly.js";
import { drainAgentDatabaseResources } from "../../state/openclaw-agent-db-resources.js";
import {
closeOpenClawAgentDatabasesAsync,
closeOpenClawAgentDatabasesForTest,
getOpenClawAgentDatabaseIfOpen,
openOpenClawAgentDatabase,
} from "../../state/openclaw-agent-db.js";
import { closeOpenClawStateDatabaseForTest } from "../../state/openclaw-state-db.js";
import * as executionCleanup from "../../state/openclaw-agent-execution-cleanup.js";
import {
closeOpenClawStateDatabaseByPathAsync,
closeOpenClawStateDatabaseForTest,
} from "../../state/openclaw-state-db.js";
import { captureOpenClawStateWorkerContext } from "../../state/openclaw-state-worker-context.js";
import { ensureSessionEntrySync } from "./session-accessor.sqlite-initial-entry.js";
import { withSqliteSessionPageReclamation } from "./session-accessor.sqlite-page-reclamation.js";
import type { SqliteReclamationWorker } from "./session-accessor.sqlite-reclamation-worker.js";
import {
createLifecycleArtifactReclamationPlan,
runSqliteSessionReclamation,
@ -33,8 +39,6 @@ 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) => {
const actual =
@ -46,18 +50,13 @@ vi.mock("./session-accessor.sqlite-reclamation-worker.js", async (importOriginal
options,
claim,
async (worker) => {
hooks.onWorker?.(worker);
const originalRun = worker.run.bind(worker);
const spy = vi.spyOn(worker, "run").mockImplementation((params) =>
originalRun({
...params,
onCommitRequest: () => {
hooks.beforeAuthorization?.();
try {
return params.onCommitRequest();
} finally {
hooks.afterAuthorization?.();
}
return params.onCommitRequest();
},
}),
);
@ -76,8 +75,6 @@ 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();
closeOpenClawStateDatabaseForTest();
@ -96,10 +93,18 @@ function createFixture() {
return { database, databaseOptions: { ...options, path: database.path } };
}
test.each(["reclaim", "worker-close"] as const)(
"reclaims pages off-thread and releases budget deferral through %s",
test.each([
"reclaim",
"resource-close",
"shared-state-close",
"resource-close-replaced",
"resource-close-removed",
] as const)(
"reclaims pages off-thread and applies current checkpoint receipts (%s)",
async (recovery) => {
const { database, databaseOptions } = createFixture();
const staleReceipt =
recovery === "resource-close-replaced" || recovery === "resource-close-removed";
const databasePathKey = sqliteReaderDatabasePathKey(database.path);
database.db.exec(`INSERT INTO cache_entries(scope, key, blob, updated_at)
VALUES ('wal-proof', 'pages', zeroblob(4194304), 1);
@ -117,28 +122,24 @@ 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 commitRequestedAtNs: bigint | undefined;
const createAdmission = workerAdmission.createSqliteWorkerOperationAdmission;
const observeAdmission = vi
.spyOn(workerAdmission, "createSqliteWorkerOperationAdmission")
.mockImplementation((admit, attachment) =>
createAdmission((request, grant) => {
if (
request.stage === "commit" &&
isRecord(request.facts) &&
isRecord(request.facts.identity) &&
typeof request.facts.identity.nativeLocation === "string" &&
sqliteReaderDatabasePathKey(request.facts.identity.nativeLocation) === databasePathKey
) {
commitRequestedAtNs = process.hrtime.bigint();
}
admit(request, grant);
}, attachment),
);
let completedAt: number | undefined;
const relayedCompletions: number[] = [];
const publish = walCheckpoint.publishSqliteWalCheckpointObservation;
@ -162,10 +163,6 @@ test.each(["reclaim", "worker-close"] as const)(
},
});
});
let retainedWorker: SqliteReclamationWorker | undefined;
hooks.onWorker = (worker) => {
retainedWorker = worker;
};
let following: Promise<void> | undefined;
try {
reader.exec("BEGIN");
@ -199,10 +196,15 @@ test.each(["reclaim", "worker-close"] as const)(
if (recovery === "reclaim") {
reader.exec("ROLLBACK");
const completed = await reclaim(7);
expect(completed.vacuumPagesRequested).toBe(7);
assert(parentReleaseNs !== undefined);
expect(completed).toMatchObject({
checkpointCompleted: true,
checkpointIncomplete: 0,
vacuumPasses: 1,
vacuumPagesRequested: 7,
});
assert(commitRequestedAtNs !== undefined);
assert(completed.checkpoint);
expect(completed.checkpoint.observedAtNs).toBeGreaterThanOrEqual(parentReleaseNs);
expect(completed.checkpoint.observedAtNs).toBeGreaterThanOrEqual(commitRequestedAtNs);
expect(original - freePages()).toBeGreaterThan(0);
expect(original - freePages()).toBeLessThanOrEqual(7);
}
@ -215,37 +217,100 @@ test.each(["reclaim", "worker-close"] as const)(
await following;
expect(ordering).toEqual(["archive-start", "archive-complete", "following-writer"]);
expect(parentReclaim).not.toHaveBeenCalled();
if (recovery === "worker-close") {
assert(retainedWorker);
if (recovery !== "reclaim") {
closeCachedOpenClawAgentDatabase(database, { eviction: true });
expect(database.db.isOpen).toBe(false);
expect(budget.checkpointBlocked).toBeDefined();
const deferredCheckpoint = budget.checkpointBlocked;
assert(deferredCheckpoint);
const observed: string[] = [];
const unsubscribe = walCheckpoint.onSqliteWalCheckpoint(({ databasePath, health }) => {
if (databasePath === databasePathKey) {
observed.push(health.state);
}
});
const cleanupFinished = createDeferredCore();
const releaseReceipt = createDeferredCore();
const cleanup = executionCleanup.cleanupRetiredAgentDatabaseLease;
const delayReceipt = staleReceipt
? vi
.spyOn(executionCleanup, "cleanupRetiredAgentDatabaseLease")
.mockImplementation(async (cleanupParams) => {
await cleanup(cleanupParams);
if (sqliteReaderDatabasePathKey(cleanupParams.lease.path) === databasePathKey) {
cleanupFinished.resolve();
await releaseReceipt.promise;
}
})
: undefined;
let closing: Promise<void> | undefined;
let closeSettled = false;
try {
reader.exec("ROLLBACK");
await retainedWorker.close();
reader.close();
expect(budget.checkpointBlocked).toBe(deferredCheckpoint);
if (recovery === "shared-state-close") {
const { admission } = captureOpenClawStateWorkerContext({ env: databaseOptions.env });
admission.assertCurrent();
closing = closeOpenClawStateDatabaseByPathAsync(admission.databasePath).then(
() => undefined,
);
expect(() => admission.assertCurrent()).toThrow("admission is closed");
} else {
closing = drainAgentDatabaseResources(
{ path: database.path, agentId: databaseOptions.agentId },
async () => undefined,
);
}
closing = closing.then(() => {
closeSettled = true;
});
if (staleReceipt) {
await Promise.race([cleanupFinished.promise, closing]);
expect(closeSettled).toBe(false);
expect(database.db.isOpen).toBe(false);
expect(reader.isOpen).toBe(false);
expect(getOpenClawAgentDatabaseIfOpen(databaseOptions)).toBeUndefined();
// The real cleanup has closed native handles and released the exact lease.
const retiredPath = `${database.path}.retired`;
fs.renameSync(database.path, retiredPath);
if (recovery === "resource-close-replaced") {
fs.copyFileSync(retiredPath, database.path);
}
releaseReceipt.resolve();
}
await closing;
expect(closeSettled).toBe(true);
const walPath = `${database.path}-wal`;
expect(fs.existsSync(walPath) ? fs.statSync(walPath).size : 0).toBe(0);
expect(observed).toContain("complete");
if (staleReceipt) {
expect(observed).not.toContain("complete");
expect(relayedCompletions).toEqual([]);
expect(budget.checkpointBlocked).toBe(deferredCheckpoint);
} else {
expect(observed).toContain("complete");
}
} finally {
releaseReceipt.resolve();
await Promise.allSettled([closing]);
delayReceipt?.mockRestore();
unsubscribe();
}
}
expect(relayedCompletions.length).toBeGreaterThan(0);
expect(budget.checkpointBlocked).toBeUndefined();
} finally {
if (reader.isTransaction) {
reader.exec("ROLLBACK");
if (!staleReceipt) {
expect(relayedCompletions.length).toBeGreaterThan(0);
expect(budget.checkpointBlocked).toBeUndefined();
}
} finally {
if (reader.isOpen) {
if (reader.isTransaction) {
reader.exec("ROLLBACK");
}
reader.close();
}
reader.close();
parentReclaim.mockRestore();
relay.mockRestore();
observeProbe.mockRestore();
observeAdmission.mockRestore();
await Promise.allSettled([following]);
}
},
);

View file

@ -39,6 +39,7 @@ import {
replaceSessionEntrySync,
} from "./session-accessor.sqlite-entry.js";
import { ensureSessionEntrySync } from "./session-accessor.sqlite-initial-entry.js";
import { withSqliteSessionPageReclamation } from "./session-accessor.sqlite-page-reclamation.js";
import {
createHistoryEvictionReclamationPlan,
createLifecycleArtifactReclamationPlan,
@ -609,28 +610,20 @@ test("one reclamation pass leaves a large freelist for bounded later maintenance
}
const budgetBefore = freePages();
const databaseOptions = plan.databaseOptions;
const duringDrain = yieldToEventLoop().then(() => {
expect(budgetBefore - freePages()).toBeGreaterThan(0);
expect(budgetBefore - freePages()).toBeLessThanOrEqual(512);
await withSqliteSessionPageReclamation(databaseOptions, async (reclaimPages) => {
const first = await reclaimPages();
expect(first.remainingFreePages).toBeGreaterThan(0);
expect(budgetBefore - first.remainingFreePages!).toBeGreaterThan(0);
expect(budgetBefore - first.remainingFreePages!).toBeLessThanOrEqual(512);
expect(database.db.isTransaction).toBe(false);
closeOpenClawAgentDatabaseByPath(database.path);
for (const scope of scopes) {
expect(appendTranscriptEventSync(scope, { type: "budget-progress" })).toEqual({
ok: true,
value: true,
});
}
await reclaimSqliteFreePages(databaseOptions, undefined, { reclaimPages });
});
// Production retains writer admission across every yielded pass. In particular,
// retiring the old handle must queue its Worker checkpoint behind this drain.
await Promise.all([
runExclusiveSqliteSessionWrite(
databaseOptions,
() => reclaimSqliteFreePages(databaseOptions),
"session.history.free-pages",
),
duringDrain,
]);
const reopened = openOpenClawAgentDatabase(databaseOptions);
expect(Number(reopened.db.prepare("PRAGMA freelist_count").get()?.freelist_count)).toBe(0);
});

View file

@ -15,6 +15,7 @@ import type {
SqliteSessionWriteDiagnostics,
} from "./session-accessor.sqlite-contract.js";
import { runExclusiveSqliteSessionWrite } from "./session-accessor.sqlite-scope.js";
import { observeSessionArchivePruning } from "./session-history-archive-pruning-diagnostics.js";
import { drainSessionStoreWriterQueuesForTest } from "./store-writer-state.test-support.js";
afterEach(() => vi.restoreAllMocks());
@ -27,9 +28,7 @@ async function readFailedWriterLog(failure: unknown, diagnostics?: SqliteSession
logging.setLoggerOverride({ level: "warn", file: logPath });
const operation = diagnostics?.artifactPreparation
? "session.lifecycle.artifacts-prepare"
: diagnostics?.archivePruning
? "session.history.archive-prune"
: "session.transcript.batch";
: "session.transcript.batch";
try {
await expect(
runExclusiveSqliteSessionWrite(
@ -155,8 +154,6 @@ test("artifact preparation file logs retain numeric phases without payload field
test("archive pruning file logs whitelist partial stage observations", async () => {
const archivePruning = {
trigger: "initial" as const,
admissionMs: 1200.4,
asyncAdmissions: 1,
checkpointCalls: 2,
checkpointIncomplete: 1,
checkpointMs: 20.6,
@ -165,13 +162,39 @@ test("archive pruning file logs whitelist partial stage observations", async ()
archiveName: "synthetic-private-archive",
content: "synthetic-private-transcript",
};
const record = await readFailedWriterLog(new Error("synthetic pruning failure"), {
archivePruning,
});
const record = await withOpenClawTestState(
{ scenario: "minimal", env: { OPENCLAW_TEST_FILE_LOG: "1" } },
async (state) => {
const logPath = state.path("archive-pruning.log");
logging.setLoggerOverride({ level: "warn", file: logPath });
const failure = new Error("synthetic pruning failure");
let clock = 0;
vi.spyOn(performance, "now").mockImplementation(() => clock);
try {
await expect(
observeSessionArchivePruning(archivePruning, async () => {
clock = 1200.4;
throw failure;
}),
).rejects.toBe(failure);
await logging.flushLogger();
const content = await fs.readFile(logPath, "utf8");
const parsed: unknown = JSON.parse(content.trim());
assert.ok(isRecord(parsed));
assert.ok(isRecord(parsed["2"]));
expect(parsed["1"]).toBe("SQLite session archive pruning failed");
expect(parsed["2"]).toHaveProperty("elapsedMs", 1200);
expect(parsed["2"]).not.toHaveProperty("queueWaitMs");
expect(parsed["2"]).not.toHaveProperty("writerExecutionMs");
return { content, details: parsed["2"] };
} finally {
await logging.flushLogger();
logging.resetLogger();
}
},
);
expect(record.details.archivePruning).toEqual({
trigger: "initial",
admissionMs: 1200,
asyncAdmissions: 1,
checkpointCalls: 2,
checkpointIncomplete: 1,
checkpointMs: 21,

View file

@ -37,7 +37,6 @@ import type {
SessionTranscriptReadScope,
SessionTranscriptWriteScope,
SqliteSessionArtifactPreparationDiagnostics,
SqliteSessionArchivePruningDiagnostics,
SqliteSessionDatabaseAdmissionDiagnostics,
SqliteSessionWriteDiagnostics,
} from "./session-accessor.sqlite-contract.js";
@ -193,39 +192,6 @@ function artifactPreparationLogFields(diagnostics: SqliteSessionArtifactPreparat
};
}
function archivePruningLogFields(diagnostics: SqliteSessionArchivePruningDiagnostics) {
const milliseconds = (value: number | undefined) =>
value === undefined ? undefined : Math.round(value);
return {
trigger: diagnostics.trigger,
admissionMs: milliseconds(diagnostics.admissionMs),
cachedAdmissions: diagnostics.cachedAdmissions,
asyncAdmissions: diagnostics.asyncAdmissions,
checkpointCalls: diagnostics.checkpointCalls,
checkpointIncomplete: diagnostics.checkpointIncomplete,
checkpoint: diagnostics.checkpoint,
totalBytesBefore: diagnostics.totalBytesBefore,
totalBytesAfter: diagnostics.totalBytesAfter,
walBytesBefore: diagnostics.walBytesBefore,
walBytesAfter: diagnostics.walBytesAfter,
checkpointMs: milliseconds(diagnostics.checkpointMs),
checkpointMaxMs: milliseconds(diagnostics.checkpointMaxMs),
vacuumMs: milliseconds(diagnostics.vacuumMs),
vacuumPasses: diagnostics.vacuumPasses,
vacuumPagesRequested: diagnostics.vacuumPagesRequested,
queryMs: milliseconds(diagnostics.queryMs),
rowDeletionMs: milliseconds(diagnostics.rowDeletionMs),
fileRemovalMs: milliseconds(diagnostics.fileRemovalMs),
removedFiles: diagnostics.removedFiles,
missingFiles: diagnostics.missingFiles,
failedRemovals: diagnostics.failedRemovals,
measurementMs: milliseconds(diagnostics.measurementMs),
measurements: diagnostics.measurements,
legacyInventoryMs: milliseconds(diagnostics.legacyInventoryMs),
completed: diagnostics.completed === true,
};
}
export async function runExclusiveSqliteSessionWrite<T>(
scope: Pick<ResolvedSqliteReadScope, "agentId" | "env" | "path">,
fn: () => Promise<T>,
@ -249,7 +215,7 @@ export async function runExclusiveSqliteSessionWrite<T>(
}
/** Observe multi-unit maintenance without retaining foreground admission between units. */
export async function observeSqliteSessionWrite<T>(
async function observeSqliteSessionWrite<T>(
scope: Pick<ResolvedSqliteReadScope, "agentId" | "env" | "path">,
fn: () => Promise<T>,
operation: SqliteSessionWriteOperation,
@ -272,9 +238,6 @@ export async function observeSqliteSessionWrite<T>(
...(diagnostics?.artifactPreparation
? { artifactPreparation: artifactPreparationLogFields(diagnostics.artifactPreparation) }
: {}),
...(diagnostics?.archivePruning
? { archivePruning: archivePruningLogFields(diagnostics.archivePruning) }
: {}),
elapsedMs: Math.round(completedAt - startedAt),
...(timing.startedAt !== undefined && timing.finishedAt !== undefined
? {

View file

@ -34,7 +34,10 @@ import { resolveMaintenanceConfigFromInput } from "./store-maintenance.js";
afterEach(() => vi.restoreAllMocks());
function observeSlowWriters(onWarning: (operation: unknown, fields: object) => void = () => {}) {
function observeSlowWriters(
onWarning: (operation: unknown, fields: object) => void = () => {},
onPruning: (fields: object) => void = () => {},
) {
let clock = 0;
// Cross the existing threshold deterministically; all storage and callback work remains real.
vi.spyOn(performance, "now").mockImplementation(() => (clock += 1_001));
@ -48,6 +51,9 @@ function observeSlowWriters(onWarning: (operation: unknown, fields: object) => v
const operation = "operation" in fields ? fields.operation : undefined;
operations.push(operation);
onWarning(operation, fields);
} else if (message === "slow SQLite session archive pruning") {
assert(fields && typeof fields === "object");
onPruning(fields);
}
return undefined;
});
@ -262,11 +268,23 @@ it("records successful archive pruning stages", async () => {
toDatabaseOptions(resolveSqliteScope({ sessionKey: historyKey, storePath })),
);
const pruning: unknown[] = [];
const operations = observeSlowWriters((operation, fields) => {
if (operation === "session.history.archive-prune" && "archivePruning" in fields) {
pruning.push(fields.archivePruning);
}
});
const aggregates: object[] = [];
const writerSegments: object[] = [];
const pageWrites: object[] = [];
observeSlowWriters(
(operation, fields) => {
if (operation === "session.history.archive-prune") {
writerSegments.push(fields);
}
if (operation === "session.history.free-pages") {
pageWrites.push(fields);
}
},
(fields) => {
aggregates.push(fields);
pruning.push("archivePruning" in fields ? fields.archivePruning : undefined);
},
);
try {
expect(
await enforceSqliteSessionHistoryDiskBudget({
@ -275,9 +293,20 @@ it("records successful archive pruning stages", async () => {
maintenance: { maxDiskBytes: 1, highWaterBytes: 0 },
}),
).toMatchObject({ removedEntries: 2 });
expect(operations).toEqual(
expect.arrayContaining(["session.history.archive-prune", "session.history.free-pages"]),
);
expect(writerSegments.length).toBeGreaterThan(0);
expect(pageWrites.length).toBeGreaterThan(0);
for (const fields of [...writerSegments, ...pageWrites]) {
expect(fields).toMatchObject({
queueWaitMs: expect.any(Number),
writerExecutionMs: expect.any(Number),
});
}
for (const fields of aggregates) {
expect(fields).toMatchObject({ elapsedMs: expect.any(Number) });
expect(fields).not.toHaveProperty("queueWaitMs");
expect(fields).not.toHaveProperty("writerExecutionMs");
expect(fields).not.toHaveProperty("completionDelayMs");
}
expect(pruning).toEqual([
expect.objectContaining({ trigger: "initial", completed: true }),
expect.objectContaining({ trigger: "after-eviction", completed: true }),
@ -286,8 +315,6 @@ it("records successful archive pruning stages", async () => {
expect(pruning[0]).toMatchObject({ checkpointIncomplete: 0 });
for (const diagnostic of pruning) {
expect(diagnostic).toMatchObject({
admissionMs: expect.any(Number),
cachedAdmissions: expect.any(Number),
checkpointCalls: expect.any(Number),
checkpointMs: expect.any(Number),
checkpointMaxMs: expect.any(Number),

View file

@ -1,30 +1,72 @@
import { performance } from "node:perf_hooks";
import { getChildLogger } from "../../logging/logger.js";
import type { SqliteSessionArchivePruningDiagnostics } from "./session-accessor.sqlite-contract.js";
function archivePruningLogFields(diagnostics: SqliteSessionArchivePruningDiagnostics) {
const milliseconds = (value: number | undefined) =>
value === undefined ? undefined : Math.round(value);
return {
trigger: diagnostics.trigger,
checkpointCalls: diagnostics.checkpointCalls,
checkpointIncomplete: diagnostics.checkpointIncomplete,
checkpoint: diagnostics.checkpoint,
totalBytesBefore: diagnostics.totalBytesBefore,
totalBytesAfter: diagnostics.totalBytesAfter,
walBytesBefore: diagnostics.walBytesBefore,
walBytesAfter: diagnostics.walBytesAfter,
checkpointMs: milliseconds(diagnostics.checkpointMs),
checkpointMaxMs: milliseconds(diagnostics.checkpointMaxMs),
vacuumMs: milliseconds(diagnostics.vacuumMs),
vacuumPasses: diagnostics.vacuumPasses,
vacuumPagesRequested: diagnostics.vacuumPagesRequested,
queryMs: milliseconds(diagnostics.queryMs),
rowDeletionMs: milliseconds(diagnostics.rowDeletionMs),
fileRemovalMs: milliseconds(diagnostics.fileRemovalMs),
removedFiles: diagnostics.removedFiles,
missingFiles: diagnostics.missingFiles,
failedRemovals: diagnostics.failedRemovals,
measurementMs: milliseconds(diagnostics.measurementMs),
measurements: diagnostics.measurements,
legacyInventoryMs: milliseconds(diagnostics.legacyInventoryMs),
completed: diagnostics.completed === true,
};
}
export async function observeSessionArchivePruning<T>(
diagnostics: SqliteSessionArchivePruningDiagnostics,
run: () => Promise<T>,
): Promise<T> {
const startedAt = performance.now();
let failed = true;
try {
const result = await run();
failed = false;
return result;
} finally {
const elapsedMs = performance.now() - startedAt;
if (failed || elapsedMs >= 1_000) {
try {
getChildLogger({ subsystem: "session-sqlite" }).warn(
failed ? "SQLite session archive pruning failed" : "slow SQLite session archive pruning",
{
elapsedMs: Math.round(elapsedMs),
archivePruning: archivePruningLogFields(diagnostics),
},
);
} catch {
// Diagnostics cannot replace the pruning result or its original failure.
}
}
}
}
type PruningStage =
| "vacuumMs"
| "queryMs"
| "rowDeletionMs"
| "fileRemovalMs"
| "measurementMs"
| "legacyInventoryMs";
export function timeArchivePruningSync<T>(
diagnostics: SqliteSessionArchivePruningDiagnostics | undefined,
stage: PruningStage,
operation: () => T,
): T {
if (!diagnostics) {
return operation();
}
const startedAt = performance.now();
try {
return operation();
} finally {
diagnostics[stage] = (diagnostics[stage] ?? 0) + performance.now() - startedAt;
}
}
export async function timeArchivePruningAsync<T>(
diagnostics: SqliteSessionArchivePruningDiagnostics | undefined,
stage: PruningStage,

View file

@ -1,64 +1,64 @@
import fs from "node:fs";
import path from "node:path";
import { setImmediate } from "node:timers/promises";
import { executeSqliteQuerySync } from "../../infra/kysely-sync.js";
import { hasErrnoCode } from "../../infra/errno.js";
import type {
SqliteWalCheckpointSnapshot,
SqliteWalHealth,
} from "../../infra/sqlite-wal-checkpoint.js";
import type { SqliteWalReclamationResult } from "../../infra/sqlite-wal-reclamation.js";
import {
openOpenClawAgentDatabase,
runOpenClawAgentWriteTransaction,
type OpenClawAgentDatabase,
type OpenClawAgentDatabaseOptions,
} from "../../state/openclaw-agent-db.js";
import type { OpenClawAgentDatabaseOptions } from "../../state/openclaw-agent-db.js";
import {
measureSessionPhysicalDiskUsage,
pruneSessionTranscriptArchivesToHighWater,
type SessionPhysicalDiskUsage,
} from "./disk-budget.js";
import type {
SqliteSessionArchivePruningDiagnostics,
SqliteSessionDatabaseAdmissionDiagnostics,
} from "./session-accessor.sqlite-contract.js";
import { getSessionKysely, withSqliteSessionDatabase } from "./session-accessor.sqlite-scope.js";
import type { SqliteSessionArchivePruningDiagnostics } from "./session-accessor.sqlite-contract.js";
import {
readSqliteSessionArchivePruning,
withSqliteSessionPageReclamation,
} from "./session-accessor.sqlite-page-reclamation.js";
import { runExclusiveSqliteSessionWrite } from "./session-accessor.sqlite-scope.js";
import {
observeSessionArchivePruning,
timeArchivePruningAsync,
timeArchivePruningSync,
} from "./session-history-archive-pruning-diagnostics.js";
async function withArchivePruningDatabase<T>(
options: OpenClawAgentDatabaseOptions,
diagnostics: SqliteSessionArchivePruningDiagnostics | undefined,
operation: (database: OpenClawAgentDatabase) => T,
): Promise<T> {
const admission: SqliteSessionDatabaseAdmissionDiagnostics | undefined = diagnostics
? {}
: undefined;
try {
return await withSqliteSessionDatabase(options, operation, undefined, admission);
} finally {
if (diagnostics && admission?.admissionMs !== undefined) {
diagnostics.admissionMs = (diagnostics.admissionMs ?? 0) + admission.admissionMs;
if (admission.admissionMode === "cached") {
diagnostics.cachedAdmissions = (diagnostics.cachedAdmissions ?? 0) + 1;
} else if (admission.admissionMode === "async") {
diagnostics.asyncAdmissions = (diagnostics.asyncAdmissions ?? 0) + 1;
}
}
}
}
import type { SessionArchivePruningOperations } from "./session-history-archive-pruning.types.js";
type PageReclamation = {
reclaimPages?: (maxPages?: number) => Promise<SqliteWalReclamationResult>;
onCheckpointIncomplete?: (checkpoint: SqliteWalCheckpointSnapshot | undefined) => void;
assertCurrent?: () => void;
};
type ArchivePruning = PageReclamation & {
withArchiveWrite: <T>(write: () => Promise<T>) => Promise<T>;
type ArchivePruningParams = Pick<PageReclamation, "onCheckpointIncomplete"> & {
archiveDirectory: string;
databaseOptions: OpenClawAgentDatabaseOptions;
diagnostics?: SqliteSessionArchivePruningDiagnostics;
highWaterBytes: number;
storePath: string;
};
type OwnedArchivePruningParams = ArchivePruningParams & {
reclaimPages: NonNullable<PageReclamation["reclaimPages"]>;
assertCurrent: () => void;
archives: SessionArchivePruningOperations;
};
function withArchivePruningWriter<T>(
params: OwnedArchivePruningParams,
run: () => Promise<T>,
): Promise<T> {
return runExclusiveSqliteSessionWrite(
params.databaseOptions,
async () => {
params.assertCurrent();
return run();
},
"session.history.archive-prune",
);
}
export type SessionArchivePruningResult = {
removedFiles: number;
usage: SessionPhysicalDiskUsage;
@ -70,8 +70,23 @@ export type SessionArchivePruningResult = {
export async function reclaimSqliteFreePages(
databaseOptions: OpenClawAgentDatabaseOptions,
diagnostics?: SqliteSessionArchivePruningDiagnostics,
limits?: PageReclamation & { maxPasses?: number; maxPages?: number; assertCurrent?: () => void },
limits?: PageReclamation & { maxPasses?: number; maxPages?: number },
): Promise<boolean> {
const reclaimPages = limits?.reclaimPages;
if (!reclaimPages) {
return withSqliteSessionPageReclamation(
databaseOptions,
(reclaim, assertCurrent, preparedOptions) =>
reclaimSqliteFreePages(preparedOptions, diagnostics, {
...limits,
reclaimPages: reclaim,
assertCurrent: () => {
limits?.assertCurrent?.();
assertCurrent();
},
}),
);
}
let remaining = limits?.maxPages;
const maxPasses = limits?.maxPasses ?? Infinity;
for (let pass = 0; pass < maxPasses && (remaining === undefined || remaining > 0); pass++) {
@ -79,12 +94,7 @@ export async function reclaimSqliteFreePages(
await setImmediate();
}
limits?.assertCurrent?.();
const result = limits?.reclaimPages
? await limits.reclaimPages(remaining)
: await withArchivePruningDatabase(databaseOptions, diagnostics, (database) => {
limits?.assertCurrent?.();
return database.walMaintenance.reclaimFreePages({ maxPages: remaining });
});
const result = await reclaimPages(remaining);
if (diagnostics) {
for (const key of [
"checkpointCalls",
@ -117,108 +127,34 @@ export async function reclaimSqliteFreePages(
return true;
}
export function hasCanonicalSessionTranscriptArchives(
export async function hasCanonicalSessionTranscriptArchives(
databaseOptions: OpenClawAgentDatabaseOptions,
): boolean {
// openclaw-agent-db.ts cache rule: LRU eviction closes idle handles across awaits.
return hasCanonicalSessionTranscriptArchivesInDatabase(
openOpenClawAgentDatabase(databaseOptions),
);
}
function hasCanonicalSessionTranscriptArchivesInDatabase(database: OpenClawAgentDatabase): boolean {
const db = getSessionKysely(database.db);
const table = executeSqliteQuerySync(
database.db,
db
.selectFrom("sqlite_schema")
.select("name")
.where("type", "=", "table")
.where("name", "=", "session_transcript_archives"),
).rows[0];
if (!table) {
return false;
}
return (
executeSqliteQuerySync(
database.db,
db
.selectFrom("session_transcript_archives")
.select("session_id")
.where("published_at", "is not", null)
.limit(1),
).rows.length > 0
);
}
function readUnpublishedSessionTranscriptArchiveNames(
database: OpenClawAgentDatabase,
): Set<string> {
const db = getSessionKysely(database.db);
const table = executeSqliteQuerySync(
database.db,
db
.selectFrom("sqlite_schema")
.select("name")
.where("type", "=", "table")
.where("name", "=", "session_transcript_archives"),
).rows[0];
if (!table) {
return new Set();
}
return new Set(
executeSqliteQuerySync(
database.db,
db
.selectFrom("session_transcript_archives")
.select("archive_name")
.where("published_at", "is", null),
).rows.map((row) => row.archive_name),
);
): Promise<boolean> {
return (await readSqliteSessionArchivePruning(databaseOptions)) !== null;
}
async function pruneCanonicalSessionTranscriptArchivesToHighWater(
params: ArchivePruning & {
archiveDirectory: string;
databaseOptions: OpenClawAgentDatabaseOptions;
diagnostics?: SqliteSessionArchivePruningDiagnostics;
highWaterBytes: number;
storePath: string;
},
params: OwnedArchivePruningParams,
): Promise<{ removedFiles: number; usage: SessionPhysicalDiskUsage }> {
const { diagnostics } = params;
let usage = await timeArchivePruningAsync(diagnostics, "measurementMs", () =>
measureSessionPhysicalDiskUsage(params.storePath),
);
const measure = () =>
timeArchivePruningAsync(diagnostics, "measurementMs", () =>
measureSessionPhysicalDiskUsage(params.storePath),
);
let usage = await measure();
let removedFiles = 0;
while (usage.totalBytes > params.highWaterBytes) {
const removed = await params.withArchiveWrite(async () => {
// Foreground work can reduce pressure while this removal waits for its FIFO turn.
usage = await timeArchivePruningAsync(diagnostics, "measurementMs", () =>
measureSessionPhysicalDiskUsage(params.storePath),
);
const removed = await withArchivePruningWriter(params, async () => {
// A foreground writer may have freed space while this item waited for admission.
usage = await measure();
params.assertCurrent();
if (usage.totalBytes <= params.highWaterBytes) {
return false;
}
const row = await withArchivePruningDatabase(
params.databaseOptions,
diagnostics,
(database) =>
timeArchivePruningSync(diagnostics, "queryMs", () => {
const db = getSessionKysely(database.db);
return executeSqliteQuerySync(
database.db,
db
.selectFrom("session_transcript_archives")
.select(["archive_name", "generation", "session_id"])
.where("published_at", "is not", null)
.orderBy("created_at", "asc")
.orderBy("session_id", "asc")
.orderBy("generation", "asc")
.limit(1),
).rows[0];
}),
const row = await timeArchivePruningAsync(diagnostics, "queryMs", () =>
params.archives.read(),
);
params.assertCurrent();
if (!row) {
return false;
}
@ -229,6 +165,7 @@ async function pruneCanonicalSessionTranscriptArchivesToHighWater(
) {
throw new Error(`Invalid canonical session archive name for ${row.session_id}`);
}
params.assertCurrent();
try {
await timeArchivePruningAsync(diagnostics, "fileRemovalMs", () =>
fs.promises.rm(archivePath),
@ -238,10 +175,8 @@ async function pruneCanonicalSessionTranscriptArchivesToHighWater(
diagnostics.removedFiles = (diagnostics.removedFiles ?? 0) + 1;
}
} catch (error) {
// SAFETY: Node filesystem failures expose the documented errno code field.
if ((error as NodeJS.ErrnoException).code !== "ENOENT") {
// The database is the recovery copy. Retain it unless its derived file
// is gone, otherwise retention could leave an undeletable orphan.
// Keep the canonical recovery copy if its derived file could not be removed.
if (!hasErrnoCode(error, "ENOENT")) {
if (diagnostics) {
diagnostics.failedRemovals = (diagnostics.failedRemovals ?? 0) + 1;
}
@ -251,33 +186,21 @@ async function pruneCanonicalSessionTranscriptArchivesToHighWater(
diagnostics.missingFiles = (diagnostics.missingFiles ?? 0) + 1;
}
}
await withArchivePruningDatabase(params.databaseOptions, diagnostics, () =>
timeArchivePruningSync(diagnostics, "rowDeletionMs", () =>
runOpenClawAgentWriteTransaction((transactionDb) => {
const transactionKysely = getSessionKysely(transactionDb.db);
executeSqliteQuerySync(
transactionDb.db,
transactionKysely
.deleteFrom("session_transcript_archives")
.where("session_id", "=", row.session_id)
.where("generation", "=", row.generation),
);
}, params.databaseOptions),
),
await timeArchivePruningAsync(diagnostics, "rowDeletionMs", () =>
params.archives.deletePublished(row),
);
return true;
});
if (!removed) {
break;
}
// Each page unit owns separate FIFO admission after this item's writer settles.
const checkpointCompleted = await reclaimSqliteFreePages(
params.databaseOptions,
diagnostics,
params,
);
usage = await timeArchivePruningAsync(diagnostics, "measurementMs", () =>
measureSessionPhysicalDiskUsage(params.storePath),
);
usage = await measure();
if (!checkpointCompleted) {
break;
}
@ -285,17 +208,30 @@ async function pruneCanonicalSessionTranscriptArchivesToHighWater(
return { removedFiles, usage };
}
export async function pruneAllSessionTranscriptArchivesToHighWater(
input: ArchivePruning & {
archiveDirectory: string;
databaseOptions: OpenClawAgentDatabaseOptions;
diagnostics?: SqliteSessionArchivePruningDiagnostics;
highWaterBytes: number;
storePath: string;
},
export function pruneAllSessionTranscriptArchivesToHighWater(
input: ArchivePruningParams,
): Promise<SessionArchivePruningResult> {
const diagnostics = input.diagnostics ?? { trigger: "initial" };
const params = { ...input, diagnostics };
return observeSessionArchivePruning(diagnostics, () =>
withSqliteSessionPageReclamation(
input.databaseOptions,
(reclaimPages, assertCurrent, databaseOptions, archives) =>
pruneSessionArchivesWithOwner({
...input,
diagnostics,
reclaimPages,
assertCurrent,
databaseOptions,
archives,
}),
),
);
}
async function pruneSessionArchivesWithOwner(
params: OwnedArchivePruningParams & { diagnostics: SqliteSessionArchivePruningDiagnostics },
): Promise<SessionArchivePruningResult> {
const { diagnostics } = params;
const measure = () =>
timeArchivePruningAsync(diagnostics, "measurementMs", () =>
measureSessionPhysicalDiskUsage(params.storePath),
@ -307,6 +243,7 @@ export async function pruneAllSessionTranscriptArchivesToHighWater(
removedFiles: number;
usage: SessionPhysicalDiskUsage;
}): SessionArchivePruningResult => {
params.assertCurrent();
diagnostics.totalBytesAfter = result.usage.totalBytes;
diagnostics.walBytesAfter = result.usage.databaseWalBytes;
const checkpointIncomplete = diagnostics.checkpointIncomplete ?? 0;
@ -322,34 +259,26 @@ export async function pruneAllSessionTranscriptArchivesToHighWater(
if (!(await reclaimSqliteFreePages(params.databaseOptions, diagnostics, params))) {
return finish({ removedFiles: 0, usage: await measure() });
}
const canonical = (await withArchivePruningDatabase(
params.databaseOptions,
diagnostics,
(database) =>
timeArchivePruningSync(diagnostics, "queryMs", () =>
hasCanonicalSessionTranscriptArchivesInDatabase(database),
),
))
? await pruneCanonicalSessionTranscriptArchivesToHighWater(params)
: { removedFiles: 0, usage: await measure() };
const canonical = await pruneCanonicalSessionTranscriptArchivesToHighWater(params);
if (diagnostics.checkpointIncomplete || canonical.usage.totalBytes <= params.highWaterBytes) {
return finish(canonical);
}
const legacy = await params.withArchiveWrite(async () =>
pruneSessionTranscriptArchivesToHighWater({
diagnostics,
excludeNames: await withArchivePruningDatabase(
params.databaseOptions,
diagnostics,
(database) =>
timeArchivePruningSync(diagnostics, "queryMs", () =>
readUnpublishedSessionTranscriptArchiveNames(database),
),
),
highWaterBytes: params.highWaterBytes,
storePath: params.storePath,
}),
);
const legacy = await pruneSessionTranscriptArchivesToHighWater({
diagnostics,
highWaterBytes: params.highWaterBytes,
storePath: params.storePath,
removeFile: (file) =>
withArchivePruningWriter(params, async () => {
const usage = await measure();
params.assertCurrent();
if (usage.totalBytes <= params.highWaterBytes) {
return "preserved";
}
return timeArchivePruningAsync(diagnostics, "fileRemovalMs", () =>
params.archives.removeLegacy(file.path),
);
}),
});
return finish({
removedFiles: canonical.removedFiles + legacy.removedFiles,
usage: legacy.usage,

View file

@ -0,0 +1,19 @@
export type PublishedSessionTranscriptArchive = {
archive_name: string;
archive_sha256: string;
created_at: number;
encoding: string;
generation: string;
published_at: number;
reason: string;
session_id: string;
session_key: string;
};
export type SessionLegacyArchiveRemovalResult = "removed" | "failed" | "preserved";
export type SessionArchivePruningOperations = {
read: () => Promise<PublishedSessionTranscriptArchive | null>;
removeLegacy: (filePath: string) => Promise<SessionLegacyArchiveRemovalResult>;
deletePublished: (archive: PublishedSessionTranscriptArchive) => Promise<void>;
};

View file

@ -0,0 +1,180 @@
import fs from "node:fs";
import path from "node:path";
import { executeSqliteQuerySync } from "../../infra/kysely-sync.js";
import type { SqliteWalReclamationResult } from "../../infra/sqlite-wal-reclamation.js";
import { assertExistingDatabaseIdentity } from "../../infra/sqlite-worker-identity.js";
import { readOpenClawAgentDatabaseIdentity } from "../../state/openclaw-agent-db-identity.js";
import {
withOpenClawAgentDatabaseReadOnly,
type OpenClawAgentReadOnlyDatabase,
} from "../../state/openclaw-agent-db-readonly.js";
import {
runOpenClawAgentWriteTransaction,
type OpenClawAgentDatabase,
type OpenClawAgentDatabaseOptions,
} from "../../state/openclaw-agent-db.js";
import { tableExists } from "../../state/openclaw-state-db-schema-helpers.js";
import { getSessionKysely } from "./session-accessor.sqlite-scope.js";
import type {
PublishedSessionTranscriptArchive,
SessionLegacyArchiveRemovalResult,
} from "./session-history-archive-pruning.types.js";
import type { SessionArchivePruningWorkerInput } from "./session-transcript-worker.types.js";
export function readSessionArchivePruningInDatabase(
database: OpenClawAgentReadOnlyDatabase,
): PublishedSessionTranscriptArchive | null {
if (!tableExists(database.db, "session_transcript_archives")) {
return null;
}
const db = getSessionKysely(database.db);
const row = executeSqliteQuerySync(
database.db,
db
.selectFrom("session_transcript_archives")
.select([
"archive_name",
"archive_sha256",
"created_at",
"encoding",
"generation",
"published_at",
"reason",
"session_id",
"session_key",
])
.where("published_at", "is not", null)
.orderBy("created_at", "asc")
.orderBy("session_id", "asc")
.orderBy("generation", "asc")
.limit(1),
).rows[0];
return row && row.published_at !== null ? { ...row, published_at: row.published_at } : null;
}
export function readSessionArchivePruningInWorker(
request: SessionArchivePruningWorkerInput,
): PublishedSessionTranscriptArchive | null {
const identity = `file:${request.expectedIdentity.physicalIdentity}`;
assertExistingDatabaseIdentity(request.database.path, identity);
const result = withOpenClawAgentDatabaseReadOnly(
(database) => {
const current = readOpenClawAgentDatabaseIdentity(database);
if (
current.identity !== request.expectedIdentity.physicalIdentity ||
current.filename !== request.expectedIdentity.nativeLocation
) {
throw new Error("SQLite archive pruning database owner changed");
}
const value = readSessionArchivePruningInDatabase(database);
assertExistingDatabaseIdentity(request.database.path, identity);
return value;
},
{ ...request.database, env: request.env },
);
if (!result.found) {
throw new Error(`SQLite archive pruning cannot read its database: ${result.reason}`);
}
return result.value;
}
function removeLegacyArchiveFile(
filePath: string,
admitCommit: () => void,
): SessionLegacyArchiveRemovalResult {
let stat: fs.Stats;
try {
stat = fs.statSync(filePath);
} catch {
admitCommit();
return "failed";
}
if (!stat.isFile()) {
admitCommit();
return "failed";
}
// Unlink cannot roll back; authorize it while the transaction still excludes peers.
admitCommit();
try {
fs.unlinkSync(filePath);
return "removed";
} catch {
return "failed";
}
}
export function removeLegacySessionArchiveInDatabase(
database: OpenClawAgentDatabase,
options: OpenClawAgentDatabaseOptions,
filePath: string,
admit: (stage: "transaction" | "commit") => void,
): SessionLegacyArchiveRemovalResult {
return runOpenClawAgentWriteTransaction((transactionDb) => {
if (transactionDb.db !== database.db) {
throw new Error("SQLite archive pruning lost its database owner");
}
admit("transaction");
const db = getSessionKysely(transactionDb.db);
const owned =
tableExists(transactionDb.db, "session_transcript_archives") &&
executeSqliteQuerySync(
transactionDb.db,
db
.selectFrom("session_transcript_archives")
.select("archive_name")
.where("archive_name", "=", path.basename(filePath))
.limit(1),
).rows.length > 0;
if (owned) {
admit("commit");
return "preserved";
}
return removeLegacyArchiveFile(filePath, () => admit("commit"));
}, options);
}
export function deletePublishedSessionArchiveInDatabase(
database: OpenClawAgentDatabase,
options: OpenClawAgentDatabaseOptions,
row: PublishedSessionTranscriptArchive,
admit: (stage: "transaction" | "commit") => void,
): void {
runOpenClawAgentWriteTransaction((transactionDb) => {
if (transactionDb.db !== database.db) {
throw new Error("SQLite archive pruning lost its database owner");
}
admit("transaction");
const db = getSessionKysely(transactionDb.db);
// A peer or cold reopen may change publication while unlink is in flight.
const deletion = executeSqliteQuerySync(
transactionDb.db,
db
.deleteFrom("session_transcript_archives")
.where("session_id", "=", row.session_id)
.where("generation", "=", row.generation)
.where("archive_name", "=", row.archive_name)
.where("archive_sha256", "=", row.archive_sha256)
.where("created_at", "=", row.created_at)
.where("encoding", "=", row.encoding)
.where("reason", "=", row.reason)
.where("session_key", "=", row.session_key)
.where("published_at", "=", row.published_at),
);
if (deletion.numAffectedRows !== 1n) {
throw new Error("SQLite session archive changed during pruning; retry cleanup.");
}
admit("commit");
}, options);
}
export function reclaimSessionArchivePagesInWorker(
database: OpenClawAgentDatabase,
maxPages: number | undefined,
admit: (stage: "transaction" | "commit") => void,
): SqliteWalReclamationResult {
return database.walMaintenance.reclaimFreePages({
maxPages,
beforeMutation: () => admit("transaction"),
onCommit: () => admit("commit"),
});
}

View file

@ -48,7 +48,6 @@ import {
} from "./session-accessor.sqlite-references.js";
import {
getSessionKysely,
observeSqliteSessionWrite,
resolveSqliteScope,
resolveSqliteTranscriptArchiveDirectory,
runExclusiveSqliteSessionWrite,
@ -120,7 +119,7 @@ export async function inspectSqliteSessionHistoryDiskBudget(
});
const databaseOptions = toDatabaseOptions(resolved);
if (
hasCanonicalSessionTranscriptArchives(databaseOptions) ||
(await hasCanonicalSessionTranscriptArchives(databaseOptions)) ||
(await hasRetainedSessionTranscriptArchives(params.storePath))
) {
return { diskBudget, wouldMutate: true };
@ -436,26 +435,15 @@ async function enforceSessionHistoryMaintenanceForDatabase(
const archiveDirectory = resolveSqliteTranscriptArchiveDirectory(resolved);
const pruneArchives = (trigger: SqliteSessionArchivePruningDiagnostics["trigger"]) => {
const archivePruning: SqliteSessionArchivePruningDiagnostics = { trigger };
return withSqliteSessionPageReclamation(databaseOptions, (reclaimPages) =>
observeSqliteSessionWrite(
resolved,
async () =>
pruneAllSessionTranscriptArchivesToHighWater({
archiveDirectory,
databaseOptions,
diagnostics: archivePruning,
highWaterBytes,
storePath: params.storePath,
reclaimPages,
withArchiveWrite: (write) =>
runExclusiveSqliteSessionWrite(resolved, write, "session.history.archive-prune"),
onCheckpointIncomplete: (checkpoint) =>
deferPhysicalBudgetForCheckpoint(params, databasePath, checkpoint),
}),
"session.history.archive-prune",
{ archivePruning },
),
);
return pruneAllSessionTranscriptArchivesToHighWater({
archiveDirectory,
databaseOptions,
diagnostics: archivePruning,
highWaterBytes,
storePath: params.storePath,
onCheckpointIncomplete: (checkpoint) =>
deferPhysicalBudgetForCheckpoint(params, databasePath, checkpoint),
});
};
let pruning = await pruneArchives("initial");
let { usage, removedFiles } = pruning;
@ -670,23 +658,20 @@ async function enforceSessionHistoryMaintenanceForDatabase(
};
const checkpointCompleted = await withSqliteSessionPageReclamation(
databaseOptions,
(reclaimPages) =>
observeSqliteSessionWrite(
resolved,
async () => {
try {
return await reclaimSqliteFreePages(databaseOptions, pageDiagnostics, {
reclaimPages,
onCheckpointIncomplete: (checkpoint) =>
deferPhysicalBudgetForCheckpoint(params, databasePath, checkpoint),
});
} catch {
// The durable deletion succeeded; a later pass can reclaim pages.
return true;
}
},
"session.history.free-pages",
),
async (reclaimPages, assertCurrent, preparedOptions) => {
try {
return await reclaimSqliteFreePages(preparedOptions, pageDiagnostics, {
reclaimPages,
assertCurrent,
onCheckpointIncomplete: (checkpoint) =>
deferPhysicalBudgetForCheckpoint(params, databasePath, checkpoint),
});
} catch {
// The durable deletion succeeded; a later pass can reclaim pages.
assertCurrent();
return true;
}
},
);
usage = await measureSessionPhysicalDiskUsage(params.storePath);
if (!checkpointCompleted) {

View file

@ -0,0 +1,129 @@
import fs from "node:fs";
import path from "node:path";
import { expect, it, vi } from "vitest";
import type { SqliteWalReclamationResult } from "../../infra/sqlite-wal-reclamation.js";
import * as tmpDirOwner from "../../infra/tmp-openclaw-dir.js";
import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js";
import { measureSessionPhysicalDiskUsage } from "./disk-budget.js";
import {
loadSessionEntryReadOnly,
patchSessionEntryCore,
} from "./session-accessor.sqlite-entry.js";
import * as pageReclamation from "./session-accessor.sqlite-page-reclamation.js";
import { loadTranscriptEventsSync } from "./session-accessor.sqlite-read.js";
import { createSessionHistoryBudgetFixture } from "./session-history-budget.test-support.js";
import { enforceSqliteSessionHistoryDiskBudget } from "./session-history-eviction.js";
it("admits a session patch between worker page-reclamation passes", async () => {
let fixtureRoot = "";
let pagesAtPatchCompletion: number | undefined;
const pageResults: SqliteWalReclamationResult[] = [];
await withOpenClawTestState(
{ prefix: "session-history-writer-fairness-", scenario: "minimal", layout: "state-only" },
async (state) => {
fixtureRoot = state.root;
const tempDir = state.sessionsDir();
fs.mkdirSync(tempDir, { recursive: true });
const storePath = path.join(tempDir, "sessions.json");
const temporaryRoot = vi
.spyOn(tmpDirOwner, "resolvePreferredOpenClawTmpDir")
.mockReturnValue(state.root);
const { createHistoricalTranscript, database, sessionExists, readArchiveNames } =
createSessionHistoryBudgetFixture(() => ({ storePath, tempDir }));
const sessionKey = "agent:main:writer-fairness";
const historical = { sessionKey, sessionId: "fairness-history", storePath };
const target = { agentId: "main", sessionKey, storePath };
let maintenance: ReturnType<typeof enforceSqliteSessionHistoryDiskBudget> | undefined;
let patch: ReturnType<typeof patchSessionEntryCore> | undefined;
const withPages = pageReclamation.withSqliteSessionPageReclamation;
const pages = vi
.spyOn(pageReclamation, "withSqliteSessionPageReclamation")
.mockImplementation(<T>(...args: Parameters<typeof withPages<T>>) => {
const [input, run] = args;
return withPages(input, (reclaim, ...context) =>
run(
async (maxPages) => {
const result = await reclaim(maxPages);
pageResults.push(result);
if (!patch) {
patch = patchSessionEntryCore(
target,
() => ({ label: "foreground writer progressed" }),
{ skipMaintenance: true, preserveActivity: true },
).then((entry) => {
pagesAtPatchCompletion = pageResults.length;
return entry;
});
void patch.catch(() => {});
}
return result;
},
...context,
),
);
});
try {
await createHistoricalTranscript({
...historical,
nextSessionId: "fairness-live",
content: "Retained history survives physical page reclamation.",
updatedAt: Date.now(),
});
const transcript = loadTranscriptEventsSync(historical);
const owner = database();
const pageSize = Number(owner.db.prepare("PRAGMA page_size").get()?.page_size);
// sqlite-allow-raw -- Synthetic cache pages exercise real incremental vacuum without deleting history.
owner.db
.prepare(
"INSERT INTO cache_entries(scope, key, blob, updated_at) VALUES (?, ?, zeroblob(?), 1)",
)
.run("writer-fairness", "free-pages", pageSize * 2048);
owner.db.prepare("DELETE FROM cache_entries WHERE scope = ?").run("writer-fairness");
owner.walMaintenance.checkpoint();
const freePages = Number(owner.db.prepare("PRAGMA freelist_count").get()?.freelist_count);
expect(freePages).toBeGreaterThan(1024);
const before = await measureSessionPhysicalDiskUsage(storePath);
const highWaterBytes = before.totalBytes - pageSize * 1024;
maintenance = enforceSqliteSessionHistoryDiskBudget({
storePath,
mode: "enforce",
maintenance: {
maxDiskBytes: before.totalBytes - 1,
highWaterBytes,
},
});
const result = await maintenance;
expect(patch).toBeDefined();
await expect(patch).resolves.toMatchObject({
sessionId: "fairness-live",
label: "foreground writer progressed",
});
expect(result).toMatchObject({ removedEntries: 0, removedFiles: 0 });
expect(result?.totalBytesAfter).toBeLessThanOrEqual(highWaterBytes);
expect(loadSessionEntryReadOnly(target)).toMatchObject({
sessionId: "fairness-live",
label: "foreground writer progressed",
});
expect(loadTranscriptEventsSync(historical)).toEqual(transcript);
expect(sessionExists("fairness-history")).toBe(true);
expect(sessionExists("fairness-live")).toBe(true);
expect(readArchiveNames("fairness-history")).toEqual([]);
expect(pageResults.length).toBeGreaterThan(1);
expect(pageResults[0]).toMatchObject({
checkpointCompleted: true,
vacuumPagesRequested: 8,
});
expect(pageResults[0]?.remainingFreePages).toBeGreaterThan(0);
expect(pageResults.at(-1)?.remainingFreePages).toBe(0);
} finally {
await Promise.allSettled([maintenance, patch]);
pages.mockRestore();
temporaryRoot.mockRestore();
}
},
);
expect(fs.existsSync(fixtureRoot)).toBe(false);
expect(pagesAtPatchCompletion).toBeGreaterThan(0);
expect(pagesAtPatchCompletion).toBeLessThan(pageResults.length);
});

View file

@ -3,18 +3,14 @@ import { channel } from "node:diagnostics_channel";
import fs from "node:fs";
import path from "node:path";
import { performance } from "node:perf_hooks";
import type { Worker, WorkerOptions } from "node:worker_threads";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { afterEach, expect, it, vi } from "vitest";
import { getNodeSqliteKysely, iterateSqliteQuerySync } from "../../infra/kysely-sync.js";
import { openNodeSqliteDatabase } from "../../infra/node-sqlite.js";
import { runtimeProcessEntrypoints } from "../../infra/runtime-process-entrypoints.js";
import { resolveRuntimeWorkerUrl } from "../../infra/runtime-worker-url.js";
import {
sqliteReaderDatabasePathKey,
withSqliteReaderOwner,
} from "../../infra/sqlite-reader-lifecycle.js";
import * as sqliteTransaction from "../../infra/sqlite-transaction.js";
import {
onSqliteWalCheckpoint,
publishSqliteWalCheckpointObservation,
@ -35,6 +31,7 @@ import {
} from "./disk-budget-runtime.js";
import type { SqliteSessionArchivePruningDiagnostics } from "./session-accessor.sqlite-contract.js";
import { runExclusiveSqliteSessionWrite } from "./session-accessor.sqlite-scope.js";
import * as archivePruningDiagnostics from "./session-history-archive-pruning-diagnostics.js";
import {
enforceSqliteSessionHistoryDiskBudget,
inspectSqliteSessionHistoryDiskBudget,
@ -42,50 +39,6 @@ import {
} from "./session-history-eviction.js";
import { resolveMaintenanceConfigFromInput } from "./store-maintenance.js";
const nativePreload = vi.hoisted(() => ({
moduleUrl: "",
path: "",
barrier: undefined as SharedArrayBuffer | undefined,
commitGate: undefined as SharedArrayBuffer | undefined,
}));
vi.mock("node:worker_threads", async (importOriginal) => {
const actual = await importOriginal<typeof import("node:worker_threads")>();
return {
...actual,
Worker: class extends actual.Worker {
private readonly selected: boolean;
constructor(filename: string | URL, options?: WorkerOptions) {
const selected = Boolean(
nativePreload.path && filename.toString() === nativePreload.moduleUrl,
);
super(
filename,
selected
? {
...options,
execArgv: [...(options?.execArgv ?? []), "--require", nativePreload.path],
workerData: { ...options?.workerData, fixtureParentBarrier: nativePreload.barrier },
}
: options,
);
this.selected = selected;
}
override postMessage(...args: Parameters<Worker["postMessage"]>) {
const request: unknown = args[0];
if (this.selected && isRecord(request) && request.type === "reclaim") {
nativePreload.commitGate =
isRecord(request.plan) &&
request.plan.kind === "maintenance-pages" &&
request.commitGate instanceof SharedArrayBuffer
? request.commitGate
: undefined;
}
return super.postMessage(...args);
}
},
};
});
const warn = vi.hoisted(() => vi.fn());
vi.mock("../../logging/subsystem.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../../logging/subsystem.js")>();
@ -101,99 +54,11 @@ vi.mock("../../logging/subsystem.js", async (importOriginal) => {
let state: OpenClawTestState;
afterEach(async () => {
vi.restoreAllMocks();
nativePreload.path = "";
nativePreload.moduleUrl = "";
nativePreload.barrier = undefined;
nativePreload.commitGate = undefined;
await drainSessionDiskBudgetWorkers();
await closeOpenClawAgentDatabasesAsync();
await state?.cleanup();
});
function installCheckpointSettlementBarrier(databasePath: string, preloadPath: string) {
const barrier = new Int32Array(new SharedArrayBuffer(3 * Int32Array.BYTES_PER_ELEMENT));
fs.writeFileSync(
preloadPath,
`
const fs = require("node:fs");
const { DatabaseSync } = require("node:sqlite");
const { workerData } = require("node:worker_threads");
const target = ${JSON.stringify(fs.realpathSync(databasePath))};
const barrier = new Int32Array(workerData.fixtureParentBarrier);
const execute = DatabaseSync.prototype.exec;
let vacuumDatabase;
DatabaseSync.prototype.exec = function(statement) {
if (Atomics.load(barrier, 0) === 1 && statement.startsWith("PRAGMA incremental_vacuum(")) {
const location = this.location();
if (location && fs.realpathSync(location) === target) vacuumDatabase = this;
}
const value = execute.call(this, statement);
if (this === vacuumDatabase && statement === "COMMIT" && Atomics.load(barrier, 1) === 0) {
Atomics.store(barrier, 1, 1);
Atomics.notify(barrier, 1);
if (Atomics.wait(barrier, 2, 0, 5_000) === "timed-out")
throw new Error("parent did not acquire the settlement barrier");
}
return value;
};
`,
);
nativePreload.path = preloadPath;
nativePreload.barrier = barrier.buffer;
nativePreload.moduleUrl = resolveRuntimeWorkerUrl(
runtimeProcessEntrypoints.sessionTranscriptArchive,
).href;
const transaction = sqliteTransaction.runSqliteImmediateTransactionSync;
let barriers = 0;
let settledInsideBarrier = false;
vi.spyOn(sqliteTransaction, "runSqliteImmediateTransactionSync").mockImplementation(
(owner, run, options) =>
transaction(
owner,
() => {
const result = run();
if (
options?.operationLabel === "session.reclamation.commit-settlement" &&
nativePreload.commitGate &&
Atomics.load(barrier, 0) === 1 &&
barriers === 0
) {
// The parent can acquire BEGIN before the worker resumes after native COMMIT.
if (Atomics.load(barrier, 1) === 0) {
expect(Atomics.wait(barrier, 1, 0, 5_000)).not.toBe("timed-out");
}
expect(Atomics.load(barrier, 1)).toBe(1);
barriers++;
const gate = new Int32Array(nativePreload.commitGate);
Atomics.store(barrier, 2, 1);
Atomics.notify(barrier, 2);
// Keep the real writer lock until the worker records native settlement.
const settled = 4;
const previous = Atomics.load(gate, 0);
if (previous !== settled) {
expect(Atomics.wait(gate, 0, previous, 5_000)).not.toBe("timed-out");
}
settledInsideBarrier = Atomics.load(gate, 0) === settled;
expect(settledInsideBarrier).toBe(true);
}
return result;
},
options,
),
);
return {
arm: () => Atomics.store(barrier, 0, 1),
verify: () => {
expect(barriers).toBe(1);
expect(settledInsideBarrier).toBe(true);
},
release: () => {
Atomics.store(barrier, 2, 1);
Atomics.notify(barrier, 2);
},
};
}
it("admits a queued foreground write before draining all vacuum batches", async () => {
state = await createOpenClawTestState({
prefix: "wal-budget-fairness-",
@ -222,11 +87,7 @@ it("admits a queued foreground write before draining all vacuum batches", async
let observedFreePages = 0;
const writes = channel("openclaw.session.write");
const observe = (message: unknown) => {
if (
!foreground &&
isRecord(message) &&
message.operation === "session.reclamation.worker-commit"
) {
if (!foreground && isRecord(message) && message.operation === "session.history.free-pages") {
foreground = runExclusiveSqliteSessionWrite(
options,
async () => {
@ -274,10 +135,6 @@ it.each(["transaction", "iterator"] as const)(
const databasePathKey = sqliteReaderDatabasePathKey(database.path);
ensureSessionTranscriptArchiveSchema(database.db);
database.db.exec("PRAGMA wal_autocheckpoint=0");
const checkpointBarrier =
kind === "transaction"
? installCheckpointSettlementBarrier(database.path, state.path("checkpoint-settlement.cjs"))
: undefined;
const sessions = state.sessionsDir();
fs.mkdirSync(sessions, { recursive: true });
const storePath = path.join(sessions, "sessions.json");
@ -351,19 +208,18 @@ it.each(["transaction", "iterator"] as const)(
reader.exec("ROLLBACK");
}
};
const diagnostics: unknown[] = [];
const writes = channel("openclaw.session.write");
const observe = (message: unknown) => {
if (
message &&
typeof message === "object" &&
"operation" in message &&
message.operation === "session.history.archive-prune"
) {
diagnostics.push(message);
}
};
writes.subscribe(observe);
const diagnostics: SqliteSessionArchivePruningDiagnostics[] = [];
const observePruning = archivePruningDiagnostics.observeSessionArchivePruning;
vi.spyOn(archivePruningDiagnostics, "observeSessionArchivePruning").mockImplementation(
async <T>(...args: Parameters<typeof observePruning<T>>) => {
const [facts, run] = args;
try {
return await observePruning(facts, run);
} finally {
diagnostics.push(structuredClone(facts));
}
},
);
warn.mockClear();
try {
const write = database.db.prepare(
@ -397,13 +253,11 @@ it.each(["transaction", "iterator"] as const)(
});
expect(diagnostics).toEqual([
expect.objectContaining({
archivePruning: expect.objectContaining({
completed: false,
checkpointCalls: 1,
checkpointIncomplete: 1,
walBytesBefore: before.databaseWalBytes,
walBytesAfter: before.databaseWalBytes,
}),
completed: false,
checkpointCalls: 1,
checkpointIncomplete: 1,
walBytesBefore: before.databaseWalBytes,
walBytesAfter: before.databaseWalBytes,
}),
]);
assert(blocked?.checkpoint);
@ -465,14 +319,8 @@ it.each(["transaction", "iterator"] as const)(
checkpointClock.mockRestore();
}
expect(fs.statSync(`${database.path}-wal`).size).toBe(0);
checkpointBarrier?.arm();
const recovered = await enforce();
checkpointBarrier?.verify();
const lastPruning = (
diagnostics.at(-1) as
| { archivePruning?: SqliteSessionArchivePruningDiagnostics }
| undefined
)?.archivePruning;
const lastPruning = diagnostics.at(-1);
expect(
recovered?.deferredReason,
recovered?.deferredReason === undefined
@ -501,14 +349,13 @@ it.each(["transaction", "iterator"] as const)(
).toBeUndefined();
expect(recovered?.totalBytesAfter).toBeLessThanOrEqual(maintenance.highWaterBytes!);
expect(diagnostics.at(-1)).toMatchObject({
archivePruning: { completed: true, checkpointIncomplete: 0 },
completed: true,
checkpointIncomplete: 0,
});
expect(() =>
JSON.stringify({ blocked, recovered, diagnostics, health: database.walMaintenance.health }),
).not.toThrow();
} finally {
checkpointBarrier?.release();
writes.unsubscribe(observe);
release();
reader.close();
}

View file

@ -21,6 +21,23 @@ export function createSessionHistoryWorkerReaders(
runRequest: SessionHistoryWorkerRequestRunner,
): Omit<SessionHistoryWorkerDatabase, "generation" | "assertCurrent"> {
return {
readArchivePruning: async (input) =>
await runRequest(
() => ({ kind: "session-archive-pruning", ...input }),
JSON.stringify(input).length * 2,
(value) => {
if (
typeof value === "boolean" ||
Array.isArray(value) ||
value.kind !== "session-archive-pruning"
) {
throw new Error(
"Session history worker returned another result instead of archive pruning",
);
}
return value.result;
},
),
readColdMetadata: async (input) =>
await runRequest(
() => ({ kind: "cold-metadata", ...input }),

View file

@ -17,6 +17,7 @@ import type {
import type { SensitiveTextRedactionSnapshot } from "../../logging/redact.js";
import type { UserTurnTranscriptAdmissionReceipt } from "../../sessions/user-turn-transcript.types.js";
import type { OpenClawRegisteredAgentDatabase } from "../../state/openclaw-agent-db-contract.js";
import type { AgentDatabaseExecutionFileIdentity } from "../../state/openclaw-agent-execution-contract.js";
import type { SessionLifecycleTimestamps } from "./lifecycle.types.js";
import type { SessionTranscriptBoundedActiveContext } from "./session-accessor.sqlite-active-context.js";
import type {
@ -51,6 +52,7 @@ import type {
} from "./session-accessor.types.js";
import type { CanonicalSessionReaderContinuation } from "./session-canonical-key.js";
import type { SessionColdArchive } from "./session-cold-storage-state.js";
import type { PublishedSessionTranscriptArchive } from "./session-history-archive-pruning.types.js";
import type {
SessionHistoryWorkerRequest,
SessionHistoryWorkerResult,
@ -348,7 +350,15 @@ export type SessionBranchSummaryWorkerInput = {
request: SessionBranchSummaryReadRequest;
};
export type SessionArchivePruningWorkerInput = {
kind: "session-archive-pruning";
database: { agentId: string; path: string };
env: NodeJS.ProcessEnv;
expectedIdentity: AgentDatabaseExecutionFileIdentity;
};
export type SessionHistoryWorkerInput =
| SessionArchivePruningWorkerInput
| SessionColdMetadataWorkerInput
| SessionTranscriptHydrationWorkerInput
| SessionTranscriptCurrentTurnEntryWorkerInput
@ -383,6 +393,10 @@ export type SessionHistoryWorkerPreparedInput = {
}[SessionHistoryDatabaseWorkerInput["kind"]];
export type SessionTranscriptWorkerValues = {
"session-archive-pruning": {
kind: "session-archive-pruning";
result: PublishedSessionTranscriptArchive | null;
};
"transcript-search": SessionTranscriptSearchWorkerResult;
"cold-metadata": SessionColdMetadataWorkerResult;
"transcript-hydration": SessionTranscriptHydrationWorkerResult;
@ -427,6 +441,9 @@ export type SessionTranscriptWorkerReply<Kind extends keyof SessionTranscriptWor
};
export type SessionHistoryWorkerDatabase = {
readArchivePruning: (
input: Omit<SessionArchivePruningWorkerInput, "kind" | "database">,
) => Promise<PublishedSessionTranscriptArchive | null>;
readColdMetadata: (
input: Omit<SessionColdMetadataWorkerInput, "kind" | "database">,
) => Promise<SessionColdMetadataWorkerResult>;

View file

@ -124,6 +124,17 @@ serveOwnedWorkerTasks(
}
}
try {
if (request.kind === "session-archive-pruning") {
const { readSessionArchivePruningInWorker } =
await import("./session-history-archive-pruning.worker.js");
return {
ok: true,
...(await withHistoryDatabase(request.database, () => ({
kind: "session-archive-pruning" as const,
result: readSessionArchivePruningInWorker(request),
}))),
};
}
if (request.kind === "cold-metadata") {
const { withOpenClawAgentDatabaseReadOnly } =
await import("../../state/openclaw-agent-db-readonly.js");

View file

@ -525,6 +525,10 @@ it.each([
const lifecycleGeneration = getAgentEventLifecycleGeneration();
const writerStarted = createDeferred();
const releaseWriter = createDeferred();
const terminalPersisted = createDeferred();
let persistenceSpy:
| MockInstance<typeof lifecycleState.persistGatewaySessionLifecycleEvent>
| undefined;
let claimId: string | undefined;
let subscriptions: ReturnType<typeof startGatewayEventSubscriptions> | undefined;
let heldWriter: Promise<unknown> | undefined;
@ -577,6 +581,21 @@ it.each([
terminalSessions: { closeTaskSessions: vi.fn() },
refreshConnectedUserProfiles: vi.fn(),
});
const persistLifecycleEvent = lifecycleState.persistGatewaySessionLifecycleEvent;
persistenceSpy = vi
.spyOn(lifecycleState, "persistGatewaySessionLifecycleEvent")
.mockImplementation((params) => {
const persistence = persistLifecycleEvent(params);
if (
params.event.runId === runId &&
params.sessionKey === target.sessionKey &&
params.event.lifecycleGeneration === lifecycleGeneration &&
params.event.data?.phase === phase
) {
terminalPersisted.resolve(persistence);
}
return persistence;
});
emitAgentEventForOwner(
{
@ -604,6 +623,8 @@ it.each([
releaseWriter.resolve();
await heldWriter;
await vi.waitFor(() => expect(loadSessionEntry(target)?.status).toBe(status));
// Failure notices join this receipt after the terminal row commits.
await terminalPersisted.promise;
await vi.waitFor(() =>
expect(getAgentRunContextOwnerStatus(runId, terminalClaimId, lifecycleGeneration)).toBe(
"clear-requested",
@ -618,6 +639,7 @@ it.each([
subscriptions?.lifecycleUnsub();
await subscriptions?.taskUnsub();
releaseAgentRunContext(runId, claimId);
persistenceSpy?.mockRestore();
routing.loadSessionEntry.mockReset();
closeOpenClawAgentDatabasesForTest();
tempDirs.cleanup();

View file

@ -12,6 +12,8 @@ import { withSqliteReaderOwner } from "./sqlite-reader-lifecycle.js";
import {
SQLITE_WORKER_MAX_RESULT_BYTES,
SQLITE_WORKER_PREPARE_COMMAND,
SQLITE_WORKER_CLOSE_RECEIPT,
type SqliteWorkerCloseReceipt,
type SqliteWorkerPreparedBackend,
type SqliteWorkerCommand,
type SqliteWorkerOperations,
@ -126,6 +128,7 @@ async function receive(request: SqliteWorkerRequest): Promise<void> {
let openNotEntered = false;
try {
let value: unknown;
let closeReceipt: SqliteWorkerCloseReceipt | undefined;
if (request.type !== "result-next" && request.type !== "execute-frame") {
if (request.lifecyclePreparation) {
const databasePath = request.stateDatabasePath ?? actorPaths.get(request.actor);
@ -480,6 +483,9 @@ async function receive(request: SqliteWorkerRequest): Promise<void> {
(SQLITE_WORKER_PREPARE_COMMAND in backend &&
backend[SQLITE_WORKER_PREPARE_COMMAND] !== undefined &&
typeof backend[SQLITE_WORKER_PREPARE_COMMAND] !== "function") ||
(SQLITE_WORKER_CLOSE_RECEIPT in backend &&
backend[SQLITE_WORKER_CLOSE_RECEIPT] !== undefined &&
typeof backend[SQLITE_WORKER_CLOSE_RECEIPT] !== "function") ||
(backend.assertSettled !== undefined && typeof backend.assertSettled !== "function") ||
(backend.prepare !== undefined && typeof backend.prepare !== "function")
) {
@ -496,6 +502,9 @@ async function receive(request: SqliteWorkerRequest): Promise<void> {
const coordinator = await prepareLifecycle();
try {
await runInActorContext(request.actor, () => backend.close());
closeReceipt = runInActorContext(request.actor, () =>
backend[SQLITE_WORKER_CLOSE_RECEIPT]?.(),
);
} catch (error) {
retire = true;
throw error;
@ -525,6 +534,7 @@ async function receive(request: SqliteWorkerRequest): Promise<void> {
id: request.id,
ok: true,
value: serialized,
...(closeReceipt ? { closeReceipt } : {}),
...(request.type === "result-next" ? { transfer: "frame" } : {}),
...(inputNext ? { input: "next" } : {}),
};

View file

@ -176,7 +176,6 @@ export function createSqliteWorkerLifecycle({
if (actor.closing) {
return actor.closing;
}
const firstAttempt = actor.cleanupState === undefined;
actor.cleanupState = "pending";
actor.closing = (async () => {
const errors: unknown[] = [];
@ -192,8 +191,6 @@ export function createSqliteWorkerLifecycle({
fail(actor.slot, error instanceof Error ? error : new Error(String(error)));
await actor.slot.exit;
}
} else if (firstAttempt && actor.slot.failed && !actor.slot.retiredAfterCompletion) {
errors.push(actor.slot.failed);
}
try {
if (

View file

@ -20,6 +20,7 @@ import {
retainSqliteWorkerErrorCode,
SqliteWorkerError,
type SqliteWorkerReply,
type SqliteWorkerCloseReceipt,
type SqliteWorkerRequest,
} from "./sqlite-worker-contract.js";
import { createSqliteWorkerLifecyclePreparation } from "./sqlite-worker-lifecycle-preparation.js";
@ -343,6 +344,7 @@ export type SqliteWorkerReplyOwner = {
error?: unknown,
value?: unknown,
settlement?: SqliteWorkerOperationSettlement,
closeReceipt?: SqliteWorkerCloseReceipt,
): void;
dispatch(): void;
};
@ -443,7 +445,11 @@ export function receiveSqliteWorkerReply(
}
settle(() => {
slot.current = undefined;
owner.finish(job, undefined, value);
if (job.request.type === "close") {
owner.finish(job, undefined, value, undefined, reply.closeReceipt);
} else {
owner.finish(job, undefined, value);
}
owner.dispatch();
});
}

View file

@ -252,7 +252,7 @@ export class SqliteWorkerBroker {
retainSqliteWorkerAdmissionCleanup(admittedActor, options.retainCleanup, () =>
this.lifecycle.closeActor(admittedActor, options.maintenanceScope),
);
options.onNativeStopped?.(actor.nativeStopped);
options.onNativeStopped?.(actor.nativeStopped, () => admittedActor.closeReceipt);
await actor.opened;
options.assertCurrent?.();
if (actor.retirementRequested) {
@ -436,7 +436,15 @@ export class SqliteWorkerBroker {
return this.lifecycle.createSlot(options, borrowedGenerationSlot, (slot) => ({
fail: (reason, currentError, completed, openOutcome) =>
this.fail(slot, reason, currentError, completed, openOutcome),
finish: (job, error, value, settlement) => this.finish(job, error, value, settlement),
finish: (job, error, value, settlement, closeReceipt) => {
if (job.request.type === "close" && closeReceipt) {
const actor = [...slot.actors].find((candidate) => candidate.id === job.request.actor);
if (actor) {
actor.closeReceipt = closeReceipt;
}
}
this.finish(job, error, value, settlement);
},
dispatch: () => this.dispatch(slot),
}));
}
@ -642,9 +650,6 @@ export class SqliteWorkerBroker {
resume(slot.failed);
}
}
if (completed) {
slot.retiredAfterCompletion = true;
}
const current = slot.current;
slot.current = undefined;
if (current) {

View file

@ -1,7 +1,11 @@
import type { Worker } from "node:worker_threads";
import type { OpenClawDatabaseMaintenanceScope } from "../state/openclaw-state-db-async-lifecycle.js";
import type { RuntimeWorkerGeneration } from "./runtime-worker-generation.js";
import type { SqliteWorkerRequest, SqliteWorkerReply } from "./sqlite-worker-contract.js";
import type {
SqliteWorkerRequest,
SqliteWorkerReply,
SqliteWorkerCloseReceipt,
} from "./sqlite-worker-contract.js";
import type {
SqliteWorkerAdmissionFactory,
SqliteWorkerOperationAdmission,
@ -66,7 +70,6 @@ export type Slot = {
queue: Job[];
current?: Job;
failed?: Error;
retiredAfterCompletion?: true;
retiring?: Promise<void>;
exit: Promise<void>;
exited: boolean;
@ -76,6 +79,7 @@ export type Actor = {
runtimeGeneration?: RuntimeWorkerGeneration;
nativeStopped: Promise<void>;
markNativeStopped(): void;
closeReceipt?: SqliteWorkerCloseReceipt;
stateDatabasePath?: string;
id: number;
key: string;
@ -144,7 +148,10 @@ export type PreparedSqliteWorkerOpen = {
createOpenAdmission?: SqliteWorkerAdmissionFactory;
maintenanceScope?: OpenClawDatabaseMaintenanceScope;
retainCleanup?: (cleanup: SqliteWorkerAdmissionCleanup) => void;
onNativeStopped?: (stopped: Promise<void>) => void;
onNativeStopped?: (
stopped: Promise<void>,
readCloseReceipt: () => SqliteWorkerCloseReceipt | undefined,
) => void;
stateDatabasePath?: string;
createAdmission?: SqliteWorkerAdmissionFactory;
assertCurrent?: () => void;

View file

@ -1,5 +1,7 @@
import type { MessagePort } from "node:worker_threads";
import type { OpenClawStateWorkerErrorPayload } from "../state/openclaw-state-worker-error.js";
import type { SqliteWalCheckpointSnapshot } from "./sqlite-wal-checkpoint.js";
import type { DatabasePathIdentity } from "./sqlite-worker-identity.js";
import type { SqliteWorkerStateContext } from "./sqlite-worker-state-context.js";
import type { SqliteWorkerTransferHandle } from "./sqlite-worker-transfer.js";
@ -17,13 +19,22 @@ export type SqliteWorkerBackend<Operations extends SqliteWorkerOperations> = {
close(): void | Promise<void>;
};
/** Recorded during successful native close; this fact never grants database access. */
export type SqliteWorkerCloseReceipt = {
identity: DatabasePathIdentity;
incarnation: string;
checkpoint: SqliteWalCheckpointSnapshot;
};
// Source fixtures and compiled backends can load separate copies in the same Worker.
export const SQLITE_WORKER_PREPARE_COMMAND = Symbol.for("openclaw.sqliteWorkerPrepareCommand");
export const SQLITE_WORKER_CLOSE_RECEIPT = Symbol.for("openclaw.sqliteWorkerCloseReceipt");
/** Internal code-loading hook; the public SDK backend remains synchronous. */
/** Internal preparation and cleanup facts; public SDK operation and close contracts stay unchanged. */
export type SqliteWorkerPreparedBackend<Operations extends SqliteWorkerOperations> =
SqliteWorkerBackend<Operations> & {
[SQLITE_WORKER_PREPARE_COMMAND]?(commandType: keyof Operations): void | Promise<void>;
[SQLITE_WORKER_CLOSE_RECEIPT]?(): SqliteWorkerCloseReceipt | undefined;
};
export type SqliteWorkerStore<Operations extends SqliteWorkerOperations> = {
@ -67,7 +78,13 @@ export type SqliteWorkerReply = {
id: number;
cleanupFailure?: OpenClawStateWorkerErrorPayload;
} & (
| { ok: true; value: Uint8Array; transfer?: "start" | "frame"; input?: "next" }
| {
ok: true;
value: Uint8Array;
transfer?: "start" | "frame";
input?: "next";
closeReceipt?: SqliteWorkerCloseReceipt;
}
| {
ok: false;
retire?: true;

View file

@ -249,7 +249,7 @@ describe("SQLite worker staged input", () => {
{ status: "rejected", reason: expect.objectContaining({ code: "unavailable" }) },
]);
stores.delete(store);
await expect(store.close()).rejects.toMatchObject({ code: "unavailable" });
await expect(store.close()).resolves.toBeUndefined();
const recovered = await open(file);
await expectRows(recovered, []);
expect(await append(recovered, value)).toMatchObject({ writes: 1 });

View file

@ -1037,7 +1037,7 @@ describe("SQLite worker store", () => {
for (const follower of followers) {
expect(follower).toMatchObject({ status: "rejected", reason: { code: "unavailable" } });
}
await expect(store.close()).rejects.toMatchObject({ code: "unavailable" });
await expect(store.close()).resolves.toBeUndefined();
stores.delete(store);
const recovered = await open(file);

View file

@ -184,7 +184,7 @@ export function openAgentDatabaseSqliteWorkerStore<Operations extends SqliteWork
custody: {
stateContext?: SqliteWorkerStateContext;
stateDatabasePath?: string;
onNativeStopped?: (stopped: Promise<void>) => void;
onNativeStopped?: SqliteWorkerOpenCustody["onNativeStopped"];
assertCurrent(): void;
createAdmission: SqliteWorkerAdmissionFactory;
},

View file

@ -3,7 +3,12 @@ import type {
SessionEntryReplacementCommit,
SessionEntryReplacementCommitted,
} from "../config/sessions/session-accessor.sqlite-replacement-state.js";
import type {
PublishedSessionTranscriptArchive,
SessionLegacyArchiveRemovalResult,
} from "../config/sessions/session-history-archive-pruning.types.js";
import type { SessionEntry } from "../config/sessions/types.js";
import type { SqliteWalReclamationResult } from "../infra/sqlite-wal-reclamation.js";
import type {
SqliteWorkerAdmissionFactory,
SqliteWorkerAdmissionRequest,
@ -42,6 +47,18 @@ export type AgentDatabaseOperations = AgentDatabaseDomainOperations & {
input: SessionProviderReviewComparison;
output: SessionEntry;
};
"session.archivePruning.deletePublished": {
input: PublishedSessionTranscriptArchive;
output: void;
};
"session.archivePruning.removeLegacy": {
input: { filePath: string };
output: SessionLegacyArchiveRemovalResult;
};
"session.archivePruning.reclaimPages": {
input: { maxPages?: number };
output: SqliteWalReclamationResult;
};
};
/** A request owner composes its retained admission with the native owner's validation. */

View file

@ -6,6 +6,8 @@ import type { Result } from "@openclaw/normalization-core/result";
import { runtimeProcessEntrypoints } from "../infra/runtime-process-entrypoints.js";
import { resolveRuntimeWorkerUrl } from "../infra/runtime-worker-url.js";
import { createSqliteLifecycleAggregateError } from "../infra/sqlite-coordinator.js";
import { publishSqliteWalCheckpointObservation } from "../infra/sqlite-wal-checkpoint.js";
import type { SqliteWorkerCloseReceipt } from "../infra/sqlite-worker-contract.js";
import { assertExistingDatabaseIdentity } from "../infra/sqlite-worker-identity.js";
import type {
SqliteWorkerAdmissionFactory,
@ -104,6 +106,7 @@ export function createAgentDatabaseNativeGeneration(
let closing: Promise<void> | undefined;
let nativeIdentity: AgentDatabaseExecutionIdentity | undefined;
let nativeStopped: Promise<void> | undefined;
let readCloseReceipt: (() => SqliteWorkerCloseReceipt | undefined) | undefined;
let lease: OpenClawAgentDatabaseWorkerLeaseReceipt | undefined;
let quickCheckPending = false;
let receiveValidation:
@ -308,8 +311,9 @@ export function createAgentDatabaseNativeGeneration(
stateDatabasePath: context.admission.databasePath,
assertCurrent,
createAdmission: admission(source, registration, assertCallerCurrent),
onNativeStopped: (stopped) => {
onNativeStopped: (stopped, readReceipt) => {
nativeStopped = stopped;
readCloseReceipt = readReceipt;
},
},
);
@ -391,6 +395,29 @@ export function createAgentDatabaseNativeGeneration(
admission(source, undefined, assertCallerCurrent),
);
}
const publishCloseCheckpoint = () => {
const receipt = readCloseReceipt?.();
if (
!receipt ||
!nativeIdentity ||
!lease ||
receipt.incarnation !== nativeIdentity.incarnation ||
receipt.identity.key !== `file:${nativeIdentity.physicalIdentity}` ||
receipt.identity.canonicalPath !== nativeIdentity.nativeLocation
) {
return;
}
try {
// Cleanup retains custody after ordinary admission is revoked during shutdown.
assertCleanupOwned();
assertExistingDatabaseIdentity(pathname, receipt.identity.key);
assertExistingDatabaseIdentity(nativeIdentity.nativeLocation, receipt.identity.key);
assertExistingDatabaseIdentity(lease.sharedStatePath, lease.sharedStateIdentity);
publishSqliteWalCheckpointObservation(pathname, receipt.checkpoint);
} catch {
// A stale diagnostic must not clear another generation's budget or fail native cleanup.
}
};
return {
failed: () =>
openingFailed || Boolean(openedStore && !isSqliteWorkerStoreAvailable(openedStore)),
@ -430,6 +457,7 @@ export function createAgentDatabaseNativeGeneration(
cause: errors[0],
});
}
publishCloseCheckpoint();
})().catch((error: unknown) => {
closing = undefined;
throw error;

View file

@ -1,8 +1,17 @@
import { MessageChannel, receiveMessageOnPort } from "node:worker_threads";
import type { Result } from "@openclaw/normalization-core/result";
import { createSqliteLifecycleAggregateError } from "../infra/sqlite-coordinator.js";
import { sqliteReaderDatabasePathKey } from "../infra/sqlite-reader-lifecycle.js";
import { assertTransactionUsable } from "../infra/sqlite-transaction.js";
import type { SqliteWorkerBackend } from "../infra/sqlite-worker-contract.js";
import {
onSqliteWalCheckpoint,
type SqliteWalCheckpointSnapshot,
} from "../infra/sqlite-wal-checkpoint.js";
import {
SQLITE_WORKER_CLOSE_RECEIPT,
type SqliteWorkerCloseReceipt,
type SqliteWorkerPreparedBackend,
} from "../infra/sqlite-worker-contract.js";
import {
assertExistingDatabaseIdentity,
readDatabasePathIdentitySync,
@ -47,7 +56,7 @@ import { openOpenClawStateDatabase } from "./openclaw-state-db.js";
export function createSqliteWorkerBackend(
input: AgentDatabaseExecutionOpen,
opening: { databasePath: string },
): SqliteWorkerBackend<AgentDatabaseOperations> {
): SqliteWorkerPreparedBackend<AgentDatabaseOperations> {
const backend = openAgentDatabaseBackend(input, opening);
try {
backend.execute({ type: "database.prepareWrite", input: undefined });
@ -70,11 +79,14 @@ export function createSqliteWorkerBackend(
export function openExistingSqliteWorkerBackend(
input: AgentDatabaseExecutionOpen,
opening: { databasePath: string; existingIdentity?: string },
): SqliteWorkerBackend<AgentDatabaseOperations> {
): SqliteWorkerPreparedBackend<AgentDatabaseOperations> {
return openAgentDatabaseBackend(input, opening);
}
type AgentDatabaseNativeBackend = Omit<SqliteWorkerBackend<AgentDatabaseOperations>, "close"> & {
type AgentDatabaseNativeBackend = Omit<
SqliteWorkerPreparedBackend<AgentDatabaseOperations>,
"close"
> & {
close(): void;
};
@ -233,6 +245,9 @@ function openAgentDatabaseBackend(
let providerReview:
| typeof import("../config/sessions/provider-review-store.worker.js")
| undefined;
let archivePruning:
| typeof import("../config/sessions/session-history-archive-pruning.worker.js")
| undefined;
let replacements:
| typeof import("../config/sessions/session-accessor.sqlite-replacement-state.js")
| undefined;
@ -247,6 +262,7 @@ function openAgentDatabaseBackend(
admit,
});
let closed = false;
let closeReceipt: SqliteWorkerCloseReceipt | undefined;
const assertOpen = () => {
if (closed) {
throw new Error("Agent database execution owner is closed");
@ -254,6 +270,17 @@ function openAgentDatabaseBackend(
};
return {
prepare(command) {
if (
command.type === "session.archivePruning.deletePublished" ||
command.type === "session.archivePruning.removeLegacy" ||
command.type === "session.archivePruning.reclaimPages"
) {
return import("../config/sessions/session-history-archive-pruning.worker.js").then(
(module) => {
archivePruning = module;
},
);
}
if (command.type === "session.entries.replace") {
return import("../config/sessions/session-accessor.sqlite-replacement-state.js").then(
(module) => {
@ -329,14 +356,60 @@ function openAgentDatabaseBackend(
admit,
);
}
if (command.type === "session.archivePruning.deletePublished" && archivePruning) {
return archivePruning.deletePublishedSessionArchiveInDatabase(
openWriter(),
options,
command.input,
admit,
);
}
if (command.type === "session.archivePruning.removeLegacy" && archivePruning) {
return archivePruning.removeLegacySessionArchiveInDatabase(
openWriter(),
options,
command.input.filePath,
admit,
);
}
if (command.type === "session.archivePruning.reclaimPages" && archivePruning) {
return archivePruning.reclaimSessionArchivePagesInWorker(
openWriter(),
command.input.maxPages,
admit,
);
}
throw new Error("Unknown agent database operation");
},
[SQLITE_WORKER_CLOSE_RECEIPT]() {
return closeReceipt;
},
close() {
closed = true;
closeReceipt = undefined;
let checkpoint: SqliteWalCheckpointSnapshot | undefined;
const errors: unknown[] = [];
for (const cleanup of [
() => domain.close(),
() => database && closeOpenClawAgentDatabaseByPath(database.path, database.agentId),
() => {
if (!database) {
return;
}
const closingPath = sqliteReaderDatabasePathKey(database.path);
const stopObserving = onSqliteWalCheckpoint((observation) => {
if (observation.databasePath === closingPath) {
checkpoint = {
health: observation.health,
observedAtNs: observation.observedAtNs,
};
}
});
try {
closeOpenClawAgentDatabaseByPath(database.path, database.agentId);
} finally {
stopObserving();
}
},
() => releaseBorrow?.(),
() => sharedBorrow?.release(),
]) {
@ -356,6 +429,16 @@ function openAgentDatabaseBackend(
errors[0],
);
}
if (identity && checkpoint) {
closeReceipt = {
identity: {
key: `file:${identity.physicalIdentity}`,
canonicalPath: identity.nativeLocation,
},
incarnation: identity.incarnation,
checkpoint,
};
}
},
};
}

View file

@ -193,6 +193,7 @@ export const databaseWorkerCoreTestFiles = [
"src/agents/sessions/session-manager-target-capture.test.ts",
"src/agents/sessions/sdk.metadata-cwd.test.ts",
"src/agents/sessions/sdk.metadata-admission.test.ts",
"src/agents/embedded-agent-runner/run/attempt-prompt-submit.retention.test.ts",
"src/agents/embedded-agent-runner/run/attempt-session-replay.test.ts",
"src/agents/embedded-agent-runner/run/attempt-stream-custody.test.ts",
"src/agents/embedded-agent-runner/run/attempt-stream-prepare.test.ts",
@ -358,6 +359,7 @@ export const databaseWorkerCoreTestFiles = [
"src/sessions/session-created.test.ts",
"src/sessions/session-upstream-links.test.ts",
"src/sessions/session-upstream-monitor.test.ts",
"src/sessions/user-turn-transcript.persistence.test.ts",
"test/canonical-descendant.integration.test.ts",
"packages/memory-host-sdk/src/host/session-memory-sync.test.ts",
"src/agents/harness/native-hook-relay-store.test.ts",
@ -373,6 +375,7 @@ export const databaseWorkerCoreTestFiles = [
"src/cli/resume-cli.test.ts",
"src/cli/update-cli/update-command-config-fence.test.ts",
"src/snapshot/git-backup.test.ts",
"src/snapshot/session-cold-storage-backup.test.ts",
"src/plugins/conversation-binding.test.ts",
"src/plugins/conversation-binding.worker.test.ts",
"src/plugins/conversation-binding.sqlite.test.ts",
@ -563,6 +566,8 @@ export const databaseWorkerCoreTestFiles = [
"src/plugin-sdk/outbound-media.retention.test.ts",
"src/plugin-sdk/provider-auth.test.ts",
"src/plugin-sdk/provider-auth-copilot-cache.test.ts",
"src/plugin-sdk/session-transcript-runtime.catalog.test.ts",
"src/plugin-sdk/session-transcript-runtime.read-fence.test.ts",
"src/plugins/doctor-contract-registry.load-paths.test.ts",
"src/skills/lifecycle/upload-store.test.ts",
"src/skills/lifecycle/upload-store-lifecycle.test.ts",

View file

@ -145,9 +145,13 @@ export const gatewayDatabaseWorkerTestFiles = [
"src/gateway/server/ws-connection/connect-hello.setup-completion.test.ts",
"src/gateway/server/ws-connection/message-handler.control-ui-build-admission.test.ts",
"src/gateway/server/ws-connection/message-handler.post-connect-health.test.ts",
"src/gateway/session-activity-summaries.test.ts",
"src/gateway/session-companion-runtime.test.ts",
"src/gateway/session-delivery-clock-jump.integration.test.ts",
"src/gateway/session-groups.test.ts",
"src/gateway/session-history-cold.benchmark.test.ts",
"src/gateway/session-history-lookup.worker.test.ts",
"src/gateway/session-history-worker-lifecycle.test.ts",
"src/gateway/session-history-worker.integration.test.ts",
"src/gateway/session-lifecycle-run-failure.test.ts",
"src/gateway/session-lifecycle-state.persistence.test.ts",
@ -162,6 +166,7 @@ export const gatewayDatabaseWorkerTestFiles = [
"src/gateway/session-startup-migration.test.ts",
"src/gateway/session-swarm-summary.test.ts",
"src/gateway/session-transcript-preview.hydration.test.ts",
"src/gateway/session-transcript-title-reader.test.ts",
"src/gateway/session-utils-store-lookup.test.ts",
"src/gateway/session-utils.agent-models.test.ts",
"src/gateway/session-utils.queued-collector-admission.test.ts",