mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
refactor(sessions): move cold maintenance selection into workers (#163200)
* refactor(sessions): move cold selection and protection to workers Keep live work admissions and cooldowns on the host while selecting and rechecking stored protection in the archive worker. Retain logical admission scopes and lexical artifact paths, and use the existing execution owner for ordinary cold maintenance opening. No schema, retention, configuration, or update behavior changes. * fix(sessions): retain freelist pragma guard annotation
This commit is contained in:
parent
b6c0584051
commit
418ec9b6be
6 changed files with 668 additions and 205 deletions
|
|
@ -158,6 +158,7 @@ const workerModules = new Set([
|
|||
"src/config/sessions/session-accessor.sqlite-mutation-worker.runtime.ts",
|
||||
"src/config/sessions/session-accessor.sqlite-summary.ts", // Only session-transcript.worker.ts dispatches the summary kernel at runtime.
|
||||
"src/config/sessions/session-accessor.sqlite-transcript-binding.ts", // History worker transcript-binding reader only.
|
||||
"src/config/sessions/session-cold-storage-selection.ts", // Cold preparation and mutation kernels in session-cold-storage-worker.ts only.
|
||||
"src/config/sessions/session-cold-storage-worker.ts", // Archive worker cold-prepare and cold-mutate dispatchers only.
|
||||
"src/config/sessions/session-membership-facts.ts", // Transcript worker session-membership-facts dispatcher only.
|
||||
|
||||
|
|
|
|||
153
src/config/sessions/session-cold-storage-selection.ts
Normal file
153
src/config/sessions/session-cold-storage-selection.ts
Normal file
|
|
@ -0,0 +1,153 @@
|
|||
import {
|
||||
executeSqliteQuerySync,
|
||||
getNodeSqliteKysely,
|
||||
iterateSqliteQuerySync,
|
||||
sqliteStringSet,
|
||||
} from "../../infra/kysely-sync.js";
|
||||
import { runSqliteDeferredTransactionSync } from "../../infra/sqlite-transaction.js";
|
||||
import { withOpenClawAgentDatabaseReadOnly } from "../../state/openclaw-agent-db-readonly.js";
|
||||
import type { DB } from "../../state/openclaw-agent-db.generated.js";
|
||||
import type {
|
||||
OpenClawAgentDatabase,
|
||||
OpenClawAgentDatabaseOptions,
|
||||
} from "../../state/openclaw-agent-db.js";
|
||||
import { readSessionStateDeleteSnapshot } from "./session-accessor.sqlite-delete-snapshot.js";
|
||||
import { collectSessionStateIdsForEntry } from "./session-accessor.sqlite-references.js";
|
||||
import { parseSessionEntryJson } from "./session-accessor.sqlite-status.js";
|
||||
import { readSessionColdStorageProtection } from "./session-cold-storage-eligibility.js";
|
||||
import { readSessionColdTranscript } from "./session-cold-storage-state.js";
|
||||
import { collectSessionAdmissionReferences } from "./session-history-eviction-candidates.js";
|
||||
import { normalizeStoreSessionKey } from "./store-entry.js";
|
||||
|
||||
export type SessionColdBatchInput = {
|
||||
databaseOptions: OpenClawAgentDatabaseOptions & { path: string };
|
||||
admissionIdentities: string[];
|
||||
cooledSessionIds: string[];
|
||||
beforeMs: number;
|
||||
maxTranscripts: number;
|
||||
maxBytes: number;
|
||||
};
|
||||
export function selectSessionColdBatch(input: SessionColdBatchInput) {
|
||||
return withOpenClawAgentDatabaseReadOnly(
|
||||
(database) =>
|
||||
runSqliteDeferredTransactionSync(
|
||||
database.db,
|
||||
() => {
|
||||
const db = getNodeSqliteKysely<DB>(database.db);
|
||||
const admissions = collectSessionAdmissionReferences({
|
||||
database,
|
||||
admissionIdentities: input.admissionIdentities,
|
||||
});
|
||||
const cooled = new Set(input.cooledSessionIds);
|
||||
const excluded = readSessionColdStorageProtection(database, input.beforeMs);
|
||||
for (const id of [...admissions, ...cooled]) {
|
||||
excluded.add(id);
|
||||
}
|
||||
const externalizations = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("session_transcript_cold_archives")
|
||||
.select("session_id")
|
||||
.where("storage", "=", "sqlite")
|
||||
.$if(cooled.size + admissions.size > 0, (query) =>
|
||||
query.where("session_id", "not in", sqliteStringSet([...cooled, ...admissions])),
|
||||
)
|
||||
.orderBy("archived_at")
|
||||
.orderBy("session_id")
|
||||
.limit(input.maxTranscripts),
|
||||
).rows.flatMap((row) => {
|
||||
const archive = readSessionColdTranscript(database.db, row.session_id);
|
||||
return archive ? [archive] : [];
|
||||
});
|
||||
const candidates =
|
||||
externalizations.length < input.maxTranscripts
|
||||
? executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("session_windows as window")
|
||||
.leftJoin(
|
||||
"session_transcript_cold_archives as cold",
|
||||
"cold.session_id",
|
||||
"window.session_id",
|
||||
)
|
||||
.select("window.session_id")
|
||||
.where("cold.session_id", "is", null)
|
||||
.$if(excluded.size > 0, (query) =>
|
||||
query.where("window.session_id", "not in", sqliteStringSet([...excluded])),
|
||||
)
|
||||
.where("window.transcript_updated_at", "<", input.beforeMs)
|
||||
.where((eb) =>
|
||||
eb.exists(
|
||||
eb
|
||||
.selectFrom("transcript_events as event")
|
||||
.select("event.seq")
|
||||
.whereRef("event.session_id", "=", "window.session_id"),
|
||||
),
|
||||
)
|
||||
.orderBy("window.transcript_updated_at")
|
||||
.orderBy("window.session_id")
|
||||
.limit(input.maxTranscripts - externalizations.length),
|
||||
).rows
|
||||
: [];
|
||||
const plans = candidates.flatMap(({ session_id: sessionId }) => {
|
||||
const snapshot = readSessionStateDeleteSnapshot(database.db, sessionId);
|
||||
return snapshot.generation && snapshot.lastSeq !== null
|
||||
? [
|
||||
{
|
||||
databaseOptions: input.databaseOptions,
|
||||
sessionId,
|
||||
beforeMs: input.beforeMs,
|
||||
snapshot,
|
||||
},
|
||||
]
|
||||
: [];
|
||||
});
|
||||
let freePages = 0;
|
||||
if (plans.length + externalizations.length === 0) {
|
||||
freePages = Number(
|
||||
// sqlite-allow-raw -- Physical maintenance is needed only when SQLite owns free pages.
|
||||
database.db.prepare("PRAGMA freelist_count").get()?.freelist_count ?? 0,
|
||||
);
|
||||
}
|
||||
return { plans, externalizations, freePages };
|
||||
},
|
||||
{ databaseLabel: database.path, operationLabel: "cold transcript selection" },
|
||||
),
|
||||
input.databaseOptions,
|
||||
);
|
||||
}
|
||||
|
||||
/** Keys whose live admissions protect any selected generation, including entry references. */
|
||||
export function readSessionAdmissionProtectionKeys(
|
||||
database: Pick<OpenClawAgentDatabase, "db">,
|
||||
sessionIds: readonly string[],
|
||||
): Set<string> {
|
||||
const selected = new Set(sessionIds);
|
||||
const keys = new Set<string>();
|
||||
if (selected.size === 0) {
|
||||
return keys;
|
||||
}
|
||||
const db = getNodeSqliteKysely<DB>(database.db);
|
||||
for (const row of iterateSqliteQuerySync(
|
||||
database.db,
|
||||
db.selectFrom("session_nodes").select(["session_key", "current_session_id", "entry_json"]),
|
||||
)) {
|
||||
const entry = parseSessionEntryJson(row);
|
||||
if (
|
||||
selected.has(row.current_session_id) ||
|
||||
(entry && collectSessionStateIdsForEntry(entry).some((id) => selected.has(id)))
|
||||
) {
|
||||
keys.add(normalizeStoreSessionKey(row.session_key));
|
||||
}
|
||||
}
|
||||
for (const row of iterateSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("session_windows")
|
||||
.select("session_key")
|
||||
.where("session_id", "in", sqliteStringSet(sessionIds)),
|
||||
)) {
|
||||
keys.add(normalizeStoreSessionKey(row.session_key));
|
||||
}
|
||||
return keys;
|
||||
}
|
||||
|
|
@ -33,6 +33,11 @@ import {
|
|||
} from "./session-cold-storage-codec.js";
|
||||
import { readSessionColdStorageProtection } from "./session-cold-storage-eligibility.js";
|
||||
import { readSessionColdStorageInventory } from "./session-cold-storage-inventory.js";
|
||||
import {
|
||||
readSessionAdmissionProtectionKeys,
|
||||
selectSessionColdBatch,
|
||||
type SessionColdBatchInput,
|
||||
} from "./session-cold-storage-selection.js";
|
||||
import {
|
||||
readSessionColdTranscript,
|
||||
type SessionColdArchive,
|
||||
|
|
@ -61,18 +66,14 @@ type SessionColdExternalization = {
|
|||
archive: Omit<SessionColdArchive, "archive_blob">;
|
||||
envelopeBytes: number;
|
||||
};
|
||||
export type SessionColdBatchInput = {
|
||||
databaseOptions: SessionColdPlan["databaseOptions"];
|
||||
plans: SessionColdPlan[];
|
||||
externalizations: Array<Omit<SessionColdArchive, "archive_blob">>;
|
||||
maxBytes: number;
|
||||
};
|
||||
export type SessionColdPreparationWorkerData = {
|
||||
type: "sqlite-transcript-archive-v2";
|
||||
operation: "cold-prepare";
|
||||
input: SessionColdBatchInput;
|
||||
};
|
||||
export type SessionColdBatchPrepared = {
|
||||
freePages: number;
|
||||
protectionKeys: string[];
|
||||
prepared: SessionColdPrepared[];
|
||||
externalizations: SessionColdExternalization[];
|
||||
oversizedSessionIds: string[];
|
||||
|
|
@ -91,6 +92,7 @@ export type SessionColdMutationPlan = { databaseOptions: SessionColdPlan["databa
|
|||
prepared: SessionColdPrepared[];
|
||||
externalizations: SessionColdExternalization[];
|
||||
beforeMs: number;
|
||||
protectionKeys: string[];
|
||||
}
|
||||
| { kind: "cold-restore"; sessionId: string; archive: Omit<SessionColdArchive, "archive_blob"> }
|
||||
);
|
||||
|
|
@ -327,15 +329,25 @@ export async function prepareSessionColdBatchInWorker(
|
|||
input: SessionColdBatchInput,
|
||||
): Promise<SessionColdBatchPrepared> {
|
||||
const result: SessionColdBatchPrepared = {
|
||||
freePages: 0,
|
||||
protectionKeys: [],
|
||||
prepared: [],
|
||||
externalizations: [],
|
||||
oversizedSessionIds: [],
|
||||
envelopeBytes: 0,
|
||||
};
|
||||
const selection = selectSessionColdBatch(input);
|
||||
if (!selection.found) {
|
||||
return result;
|
||||
}
|
||||
result.freePages = selection.value.freePages;
|
||||
const maxBytes = Math.min(input.maxBytes, MAX_COLD_ARCHIVE_BYTES);
|
||||
const candidates = [
|
||||
...input.externalizations.map((archive) => ({ kind: "externalize" as const, archive })),
|
||||
...input.plans.map((plan) => ({ kind: "archive" as const, plan })),
|
||||
...selection.value.externalizations.map((archive) => ({
|
||||
kind: "externalize" as const,
|
||||
archive,
|
||||
})),
|
||||
...selection.value.plans.map((plan) => ({ kind: "archive" as const, plan })),
|
||||
];
|
||||
for (const candidate of candidates.slice(0, 128)) {
|
||||
const remaining = maxBytes - result.envelopeBytes;
|
||||
|
|
@ -368,6 +380,26 @@ export async function prepareSessionColdBatchInWorker(
|
|||
break;
|
||||
}
|
||||
}
|
||||
const included = [
|
||||
...result.prepared.map((item) => item.plan.sessionId),
|
||||
...result.externalizations.map((item) => item.archive.session_id),
|
||||
];
|
||||
if (included.length > 0) {
|
||||
const protection = withOpenClawAgentDatabaseReadOnly(
|
||||
(database) =>
|
||||
runSqliteDeferredTransactionSync(
|
||||
database.db,
|
||||
() => [...readSessionAdmissionProtectionKeys(database, included)],
|
||||
{ databaseLabel: database.path, operationLabel: "cold admission protection" },
|
||||
),
|
||||
input.databaseOptions,
|
||||
);
|
||||
if (!protection.found) {
|
||||
throw new Error("Cold transcript database disappeared during preparation");
|
||||
}
|
||||
result.protectionKeys = protection.value;
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
|
|
@ -459,6 +491,13 @@ export function mutateSessionColdTranscriptInWorker(
|
|||
return result;
|
||||
}
|
||||
if (plan.kind === "cold-batch") {
|
||||
const protectionKeys = readSessionAdmissionProtectionKeys(database, [
|
||||
...plan.prepared.map((item) => item.plan.sessionId),
|
||||
...plan.externalizations.map((item) => item.archive.session_id),
|
||||
]);
|
||||
if ([...protectionKeys].some((key) => !plan.protectionKeys.includes(key))) {
|
||||
throw new Error("Transcript ownership changed; cold archival was canceled");
|
||||
}
|
||||
const protectedIds = readSessionColdStorageProtection(database, plan.beforeMs);
|
||||
const archivedIds: string[] = [];
|
||||
for (const prepared of plan.prepared) {
|
||||
|
|
|
|||
250
src/config/sessions/session-cold-storage.selection.test.ts
Normal file
250
src/config/sessions/session-cold-storage.selection.test.ts
Normal file
|
|
@ -0,0 +1,250 @@
|
|||
import fs from "node:fs/promises";
|
||||
import { afterAll, afterEach, beforeAll, beforeEach, expect, it, vi } from "vitest";
|
||||
import { awaitGateBeforeSettlement, createDeferred } from "../../../test/helpers/promise.js";
|
||||
import {
|
||||
emptySqliteCounts,
|
||||
observeParentSqlite,
|
||||
} from "../../../test/helpers/sqlite-parent-observer.js";
|
||||
import { beginSessionWorkAdmission } from "../../sessions/session-lifecycle-admission.js";
|
||||
import { closeOpenClawAgentDatabasesAsync } from "../../state/openclaw-agent-db.js";
|
||||
import { closeOpenClawStateDatabaseAsync } from "../../state/openclaw-state-db.js";
|
||||
import {
|
||||
createOpenClawTestState,
|
||||
type OpenClawTestState,
|
||||
} from "../../test-utils/openclaw-test-state.js";
|
||||
import * as archiveWorkers from "./session-accessor.sqlite-archive.js";
|
||||
import { replaceSessionEntrySync } from "./session-accessor.sqlite-entry.js";
|
||||
import { resolveSessionColdArchivePath } from "./session-cold-storage-codec.js";
|
||||
import { readSessionColdTranscript } from "./session-cold-storage-state.js";
|
||||
import { runSessionColdStorageMaintenance } from "./session-cold-storage.js";
|
||||
import {
|
||||
createSessionColdStorageFixture,
|
||||
historicalId,
|
||||
maintenanceConfig,
|
||||
} from "./session-cold-storage.test-support.js";
|
||||
|
||||
let state: OpenClawTestState;
|
||||
let fixture: Awaited<ReturnType<typeof createSessionColdStorageFixture>>;
|
||||
let seedPath: string;
|
||||
let aliasStorePath: string;
|
||||
|
||||
beforeAll(async () => {
|
||||
state = await createOpenClawTestState({ scenario: "minimal" });
|
||||
fixture = await createSessionColdStorageFixture(state.statePath("selection.sqlite"));
|
||||
await closeOpenClawAgentDatabasesAsync();
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
seedPath = state.statePath("selection-seed.sqlite");
|
||||
await fs.copyFile(fixture.scope.storePath, seedPath);
|
||||
aliasStorePath = state.statePath("alias", "selection.sqlite");
|
||||
if (process.platform === "win32") {
|
||||
await fs.symlink(state.stateDir, state.statePath("alias"), "junction");
|
||||
} else {
|
||||
await fs.mkdir(state.statePath("alias"));
|
||||
await fs.symlink(fixture.scope.storePath, aliasStorePath, "file");
|
||||
}
|
||||
});
|
||||
beforeEach(async () => {
|
||||
await fs.copyFile(seedPath, fixture.scope.storePath);
|
||||
});
|
||||
afterEach(async () => {
|
||||
vi.restoreAllMocks();
|
||||
await closeOpenClawAgentDatabasesAsync();
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
});
|
||||
afterAll(async () => {
|
||||
await state.cleanup();
|
||||
});
|
||||
|
||||
function delayPreparation() {
|
||||
const entered = createDeferred();
|
||||
const release = createDeferred();
|
||||
const run = archiveWorkers.runSqliteTranscriptArchiveWorkerOperation;
|
||||
const worker = vi
|
||||
.spyOn(archiveWorkers, "runSqliteTranscriptArchiveWorkerOperation")
|
||||
.mockImplementation(async (params) => {
|
||||
const result = await run(params);
|
||||
if (
|
||||
params.expectedMessageType === "done" &&
|
||||
"operation" in params.workerData &&
|
||||
params.workerData.operation === "cold-prepare"
|
||||
) {
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
}
|
||||
return result;
|
||||
});
|
||||
return { entered, release, worker };
|
||||
}
|
||||
|
||||
function expectHistoryUnchanged() {
|
||||
expect(fixture.snapshot()).toEqual(fixture.original);
|
||||
expect(readSessionColdTranscript(fixture.database(), historicalId)).toBeUndefined();
|
||||
}
|
||||
|
||||
it.each([
|
||||
{ name: "a newly admitted normalized logical key", protectsHistory: true },
|
||||
{ name: "an unrelated newly admitted key", protectsHistory: false },
|
||||
])("rechecks $name after worker selection without parent SQLite", async ({ protectsHistory }) => {
|
||||
const ownerStorePath = protectsHistory ? state.statePath("selection.json") : aliasStorePath;
|
||||
const delayed = delayPreparation();
|
||||
const observer = observeParentSqlite();
|
||||
const pending = runSessionColdStorageMaintenance({
|
||||
config: maintenanceConfig(ownerStorePath),
|
||||
});
|
||||
const outcome = pending.catch((error: unknown) => error);
|
||||
let admission: Awaited<ReturnType<typeof beginSessionWorkAdmission>> | undefined;
|
||||
try {
|
||||
await awaitGateBeforeSettlement(
|
||||
delayed.entered.promise,
|
||||
pending,
|
||||
"Cold selection was not dispatched",
|
||||
);
|
||||
admission = await beginSessionWorkAdmission({
|
||||
scope: ownerStorePath,
|
||||
identities: [
|
||||
protectsHistory ? fixture.scope.sessionKey.toUpperCase() : "agent:main:unrelated-work",
|
||||
],
|
||||
assertAllowed: () => {},
|
||||
});
|
||||
delayed.release.resolve();
|
||||
if (protectsHistory) {
|
||||
await expect(pending).rejects.toThrow("Transcript became active");
|
||||
expect(
|
||||
delayed.worker.mock.calls.some(([params]) => params.expectedMessageType === "reclaimed"),
|
||||
).toBe(false);
|
||||
} else {
|
||||
await expect(pending).resolves.toEqual({
|
||||
archivedTranscripts: 1,
|
||||
externalizedTranscripts: 0,
|
||||
});
|
||||
}
|
||||
expect(observer.counts).toEqual(emptySqliteCounts());
|
||||
} finally {
|
||||
delayed.release.resolve();
|
||||
await outcome;
|
||||
observer.restore();
|
||||
admission?.release();
|
||||
await admission?.released;
|
||||
}
|
||||
if (protectsHistory) {
|
||||
expectHistoryUnchanged();
|
||||
} else {
|
||||
const archive = readSessionColdTranscript(fixture.database(), historicalId)!;
|
||||
expect(
|
||||
(await fs.stat(resolveSessionColdArchivePath(ownerStorePath, archive.archive_name))).size,
|
||||
).toBe(archive.archive_bytes);
|
||||
}
|
||||
});
|
||||
|
||||
it("refuses revoked configuration after preparation before dispatching a mutation", async () => {
|
||||
const delayed = delayPreparation();
|
||||
const config = maintenanceConfig(fixture.scope.storePath);
|
||||
const pending = runSessionColdStorageMaintenance({
|
||||
config,
|
||||
assertCurrent: () => {
|
||||
if (!config.session.maintenance.coldStorage.enabled) {
|
||||
throw new Error("Cold maintenance configuration was revoked");
|
||||
}
|
||||
},
|
||||
});
|
||||
const outcome = pending.catch((error: unknown) => error);
|
||||
try {
|
||||
await awaitGateBeforeSettlement(
|
||||
delayed.entered.promise,
|
||||
pending,
|
||||
"Cold selection was not dispatched",
|
||||
);
|
||||
config.session.maintenance.coldStorage.enabled = false;
|
||||
delayed.release.resolve();
|
||||
await expect(pending).rejects.toThrow("Cold maintenance configuration was revoked");
|
||||
expect(
|
||||
delayed.worker.mock.calls.some(([params]) => params.expectedMessageType === "reclaimed"),
|
||||
).toBe(false);
|
||||
} finally {
|
||||
delayed.release.resolve();
|
||||
await outcome;
|
||||
}
|
||||
expectHistoryUnchanged();
|
||||
});
|
||||
|
||||
it("propagates selection failure without a mutation or synchronous fallback", async () => {
|
||||
const failure = new Error("Cold selection worker refused");
|
||||
const worker = vi
|
||||
.spyOn(archiveWorkers, "runSqliteTranscriptArchiveWorkerOperation")
|
||||
.mockRejectedValueOnce(failure);
|
||||
const observer = observeParentSqlite();
|
||||
try {
|
||||
await expect(
|
||||
runSessionColdStorageMaintenance({ config: maintenanceConfig(fixture.scope.storePath) }),
|
||||
).rejects.toBe(failure);
|
||||
expect(observer.counts).toEqual(emptySqliteCounts());
|
||||
expect(worker).toHaveBeenCalledOnce();
|
||||
} finally {
|
||||
observer.restore();
|
||||
}
|
||||
expectHistoryUnchanged();
|
||||
});
|
||||
|
||||
it("retains an embedded archive when its window gains an admitted owner after preparation", async () => {
|
||||
const config = maintenanceConfig(fixture.scope.storePath);
|
||||
await expect(runSessionColdStorageMaintenance({ config })).resolves.toEqual({
|
||||
archivedTranscripts: 1,
|
||||
externalizedTranscripts: 0,
|
||||
});
|
||||
const archive = readSessionColdTranscript(fixture.database(), historicalId)!;
|
||||
const bytes = await fs.readFile(
|
||||
resolveSessionColdArchivePath(fixture.scope.storePath, archive.archive_name),
|
||||
);
|
||||
const reboundKey = "agent:main:rebound-cold-owner";
|
||||
const reboundScope = {
|
||||
...fixture.scope,
|
||||
sessionKey: reboundKey,
|
||||
sessionId: "rebound-cold-current",
|
||||
};
|
||||
replaceSessionEntrySync(reboundScope, { sessionId: reboundScope.sessionId, updatedAt: 1 });
|
||||
fixture
|
||||
.database()
|
||||
.prepare(
|
||||
"UPDATE session_transcript_cold_archives SET storage = 'sqlite', archive_blob = ? WHERE session_id = ?",
|
||||
)
|
||||
.run(bytes, historicalId);
|
||||
await closeOpenClawAgentDatabasesAsync();
|
||||
|
||||
const delayed = delayPreparation();
|
||||
const pending = runSessionColdStorageMaintenance({ config });
|
||||
const outcome = pending.catch((error: unknown) => error);
|
||||
let admission: Awaited<ReturnType<typeof beginSessionWorkAdmission>> | undefined;
|
||||
try {
|
||||
await awaitGateBeforeSettlement(
|
||||
delayed.entered.promise,
|
||||
pending,
|
||||
"Cold externalization was not prepared",
|
||||
);
|
||||
fixture
|
||||
.database()
|
||||
.prepare("UPDATE session_windows SET session_key = ? WHERE session_id = ?")
|
||||
.run(reboundKey, historicalId);
|
||||
const before = fixture.snapshot();
|
||||
admission = await beginSessionWorkAdmission({
|
||||
scope: fixture.scope.storePath,
|
||||
identities: [reboundKey],
|
||||
assertAllowed: () => {},
|
||||
});
|
||||
delayed.release.resolve();
|
||||
await expect(pending).rejects.toThrow("Transcript ownership changed");
|
||||
expect(fixture.snapshot()).toEqual(before);
|
||||
expect(
|
||||
fixture
|
||||
.database()
|
||||
.prepare(
|
||||
"SELECT storage, archive_blob FROM session_transcript_cold_archives WHERE session_id = ?",
|
||||
)
|
||||
.get(historicalId),
|
||||
).toEqual({ storage: "sqlite", archive_blob: new Uint8Array(bytes) });
|
||||
} finally {
|
||||
delayed.release.resolve();
|
||||
await outcome;
|
||||
admission?.release();
|
||||
await admission?.released;
|
||||
}
|
||||
});
|
||||
|
|
@ -26,6 +26,7 @@ import { clearOpenClawAgentIntegrityVerification } from "../../state/openclaw-qu
|
|||
import { closeOpenClawStateDatabaseAsync } from "../../state/openclaw-state-db.js";
|
||||
import { replaceSessionEntry } from "./session-accessor.js";
|
||||
import * as archiveWorkers from "./session-accessor.sqlite-archive.js";
|
||||
import type { SqliteSessionReclamationDiagnostics } from "./session-accessor.sqlite-contract.js";
|
||||
import { readSessionTranscriptHistoryEvents } from "./session-accessor.sqlite-history.test-support.js";
|
||||
import { planSessionStateDeleteIfUnreferenced } from "./session-accessor.sqlite-lifecycle-state.js";
|
||||
import {
|
||||
|
|
@ -295,18 +296,22 @@ describe("cold transcript storage workers", () => {
|
|||
let clock = 0;
|
||||
vi.spyOn(performance, "now").mockImplementation(() => clock);
|
||||
const originalWorker = archiveWorkers.runSqliteTranscriptArchiveWorkerOperation;
|
||||
const archiveDiagnostics: SqliteSessionReclamationDiagnostics[] = [];
|
||||
vi.spyOn(archiveWorkers, "runSqliteTranscriptArchiveWorkerOperation").mockImplementation(
|
||||
(params) => {
|
||||
if (params.expectedMessageType !== "reclaimed") {
|
||||
return originalWorker(params);
|
||||
}
|
||||
if (params.diagnostics) {
|
||||
archiveDiagnostics.push(params.diagnostics);
|
||||
}
|
||||
const withWriteAdmission = params.withWriteAdmission;
|
||||
let admissions = 0;
|
||||
return originalWorker({
|
||||
...params,
|
||||
withWriteAdmission: async (...args) => {
|
||||
// Measure waiting between actual Worker admissions without delaying native work.
|
||||
if (++admissions === 2) {
|
||||
// Measure time waiting for actual Worker admission without delaying native work.
|
||||
if (++admissions === 1) {
|
||||
clock += 1_500;
|
||||
}
|
||||
return withWriteAdmission(...args);
|
||||
|
|
@ -314,8 +319,8 @@ describe("cold transcript storage workers", () => {
|
|||
});
|
||||
},
|
||||
);
|
||||
const workers: Worker[] = [];
|
||||
const observeWorker = (worker: Worker) => workers.push(worker);
|
||||
const workers = new Map<number, Worker>();
|
||||
const observeWorker = (worker: Worker) => workers.set(worker.threadId, worker);
|
||||
process.on("worker", observeWorker);
|
||||
try {
|
||||
closeOpenClawAgentDatabasesForTest();
|
||||
|
|
@ -357,8 +362,10 @@ describe("cold transcript storage workers", () => {
|
|||
exitCode: 0,
|
||||
})),
|
||||
);
|
||||
expect(workers.length).toBeGreaterThan(0);
|
||||
expect(workers.every((worker) => worker.threadId === -1)).toBe(true);
|
||||
expect(archiveDiagnostics).toHaveLength(3);
|
||||
for (const diagnostics of archiveDiagnostics) {
|
||||
expect(workers.get(diagnostics.workerThreadId!)?.threadId).toBe(-1);
|
||||
}
|
||||
} finally {
|
||||
process.off("worker", observeWorker);
|
||||
vi.restoreAllMocks();
|
||||
|
|
@ -485,13 +492,14 @@ describe("cold transcript storage workers", () => {
|
|||
).toBeUndefined();
|
||||
});
|
||||
|
||||
it.each(["authority", "retained claim"])(
|
||||
it.each(["authority", "database owner"])(
|
||||
"rolls back every candidate when maintenance %s is revoked at commit",
|
||||
async (revocation) => {
|
||||
const fixture = await createBatchFixture();
|
||||
await closeOpenClawAgentDatabaseByPathAsync(fixture.options.path);
|
||||
fixture.database();
|
||||
const originalWorker = archiveWorkers.runSqliteTranscriptArchiveWorkerOperation;
|
||||
const archiveDiagnostics: SqliteSessionReclamationDiagnostics[] = [];
|
||||
let prepared = false;
|
||||
let revoked = false;
|
||||
let closeElapsedMs = 0;
|
||||
|
|
@ -507,12 +515,15 @@ describe("cold transcript storage workers", () => {
|
|||
return result;
|
||||
}
|
||||
if (prepared && params.expectedMessageType === "reclaimed") {
|
||||
if (params.diagnostics) {
|
||||
archiveDiagnostics.push(params.diagnostics);
|
||||
}
|
||||
return originalWorker({
|
||||
...params,
|
||||
onCommitRequest: () => {
|
||||
const started = performance.now();
|
||||
revoked =
|
||||
revocation === "retained claim"
|
||||
revocation === "database owner"
|
||||
? closeOpenClawAgentDatabaseByPath(fixture.options.path)
|
||||
: true;
|
||||
closeElapsedMs = performance.now() - started;
|
||||
|
|
@ -523,8 +534,8 @@ describe("cold transcript storage workers", () => {
|
|||
return originalWorker(params);
|
||||
},
|
||||
);
|
||||
const workers: Worker[] = [];
|
||||
const observeWorker = (worker: Worker) => workers.push(worker);
|
||||
const workers = new Map<number, Worker>();
|
||||
const observeWorker = (worker: Worker) => workers.set(worker.threadId, worker);
|
||||
process.on("worker", observeWorker);
|
||||
try {
|
||||
await expect(
|
||||
|
|
@ -539,10 +550,10 @@ describe("cold transcript storage workers", () => {
|
|||
).rejects.toThrow(
|
||||
revocation === "authority"
|
||||
? "Test maintenance authority was revoked"
|
||||
: "claim is no longer current",
|
||||
: "Agent database execution admission is closed",
|
||||
);
|
||||
expect(workers.length).toBeGreaterThan(0);
|
||||
expect(workers.every((worker) => worker.threadId === -1)).toBe(true);
|
||||
expect(archiveDiagnostics).toHaveLength(1);
|
||||
expect(workers.get(archiveDiagnostics[0]!.workerThreadId!)?.threadId).toBe(-1);
|
||||
} finally {
|
||||
process.off("worker", observeWorker);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,20 +1,16 @@
|
|||
import { statSync } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { hasErrnoCode } from "../../infra/errno.js";
|
||||
import {
|
||||
executeSqliteQuerySync,
|
||||
getNodeSqliteKysely,
|
||||
sqliteStringSet,
|
||||
} from "../../infra/kysely-sync.js";
|
||||
import { createSubsystemLogger } from "../../logging/subsystem.js";
|
||||
import { KeyedAsyncQueue } from "../../plugin-sdk/keyed-async-queue.js";
|
||||
import { isIncognitoSessionKey } from "../../routing/session-key.js";
|
||||
import { collectActiveSessionWorkAdmissions } from "../../sessions/session-lifecycle-admission.js";
|
||||
import { sessionChanges } from "../../sessions/session-row-changes.js";
|
||||
import {
|
||||
retainOpenClawAgentDatabaseReadOnly,
|
||||
withOpenClawAgentDatabaseReadOnly,
|
||||
} from "../../state/openclaw-agent-db-readonly.js";
|
||||
import type { DB } from "../../state/openclaw-agent-db.generated.js";
|
||||
import { prepareOpenClawAgentDatabaseRegistrySnapshotRead } from "../../state/openclaw-agent-db-registry-listing.js";
|
||||
import {
|
||||
resolveOpenClawAgentSqlitePath,
|
||||
type OpenClawAgentDatabaseOptions,
|
||||
|
|
@ -23,6 +19,10 @@ import {
|
|||
createOpenClawAgentDatabasePathMatcher,
|
||||
isIncognitoOpenClawAgentSqlitePath,
|
||||
} from "../../state/openclaw-agent-db.paths.js";
|
||||
import {
|
||||
captureOpenClawAgentDatabaseExecution,
|
||||
supportsOpenClawAgentDatabaseExecution,
|
||||
} from "../../state/openclaw-agent-execution.js";
|
||||
import {
|
||||
resolveOpenClawStateDirForDatabasePath,
|
||||
resolveOpenClawStateSqlitePath,
|
||||
|
|
@ -34,9 +34,9 @@ import type {
|
|||
SessionTranscriptReadScope,
|
||||
SqliteSessionReclamationDiagnostics,
|
||||
} from "./session-accessor.sqlite-contract.js";
|
||||
import { readSessionStateDeleteSnapshot } from "./session-accessor.sqlite-delete-snapshot.js";
|
||||
import { withSqliteSessionPageReclamation } from "./session-accessor.sqlite-page-reclamation.js";
|
||||
import { withSqliteReclamationAuthorization } from "./session-accessor.sqlite-reclamation-commit.js";
|
||||
import { withSessionEntryWorker } from "./session-accessor.sqlite-replacement-worker.js";
|
||||
import {
|
||||
prepareSqliteTranscriptReadScope,
|
||||
resolveSqliteTranscriptReadScope,
|
||||
|
|
@ -44,21 +44,21 @@ import {
|
|||
toDatabaseOptions,
|
||||
} from "./session-accessor.sqlite-scope.js";
|
||||
import { withSqliteMutationWorkerLifetime } from "./session-accessor.sqlite-worker-request.js";
|
||||
import { readSessionColdStorageProtection } from "./session-cold-storage-eligibility.js";
|
||||
import type { SessionColdReadPreparation } from "./session-cold-storage-read.js";
|
||||
import { readSessionColdTranscript } from "./session-cold-storage-state.js";
|
||||
import type {
|
||||
SessionColdMutationPlan,
|
||||
SessionColdBatchInput,
|
||||
SessionColdBatchPrepared,
|
||||
SessionColdMutationResult,
|
||||
SessionColdPreparationWorkerData,
|
||||
SessionColdWorkerData,
|
||||
} from "./session-cold-storage-worker.js";
|
||||
import { reclaimSqliteFreePages } from "./session-history-archive-pruning.js";
|
||||
import { collectAdmissionProtectedSessionIds } from "./session-history-eviction.js";
|
||||
import { prepareSessionStoreTargetInventory } from "./session-store-target-inventory.js";
|
||||
import { withSessionHistoryWorkerReadCandidates } from "./session-transcript-worker-resources.js";
|
||||
import { withSessionHistoryWorkerDatabase } from "./session-transcript-worker-runtime.js";
|
||||
import { resolveSessionStoreTargets } from "./targets.js";
|
||||
import { normalizeStoreSessionKey } from "./store-entry.js";
|
||||
import { listConfiguredSessionStoreAgentIds } from "./targets.js";
|
||||
import { captureSessionTranscriptStorageEnvironment } from "./transcript-target-binding.js";
|
||||
|
||||
const operations = new KeyedAsyncQueue();
|
||||
|
|
@ -95,20 +95,42 @@ async function runColdMutation(
|
|||
return await withSqliteMutationWorkerLifetime(
|
||||
plan.databaseOptions,
|
||||
async ({ assertCurrent: assertRequestCurrent, commitGate, signal }) => {
|
||||
const retained = await runExclusiveSqliteSessionWrite(
|
||||
plan.databaseOptions,
|
||||
async () => {
|
||||
const execution =
|
||||
plan.kind !== "cold-restore" && supportsOpenClawAgentDatabaseExecution(plan.databaseOptions)
|
||||
? captureOpenClawAgentDatabaseExecution(plan.databaseOptions)
|
||||
: undefined;
|
||||
let retained: ReturnType<typeof retainOpenClawAgentDatabaseReadOnly> | undefined;
|
||||
try {
|
||||
const assertOpeningCurrent = () => {
|
||||
assertRequestCurrent();
|
||||
assertCurrent?.();
|
||||
return retainOpenClawAgentDatabaseReadOnly(plan.databaseOptions);
|
||||
},
|
||||
"session.reclamation.retain",
|
||||
);
|
||||
if (!retained.found) {
|
||||
throw new Error("Cold transcript operation lost its owning database");
|
||||
}
|
||||
const { database, claim } = retained;
|
||||
try {
|
||||
};
|
||||
if (!execution) {
|
||||
retained = await runExclusiveSqliteSessionWrite(
|
||||
plan.databaseOptions,
|
||||
async () => {
|
||||
assertOpeningCurrent();
|
||||
return retainOpenClawAgentDatabaseReadOnly(plan.databaseOptions);
|
||||
},
|
||||
"session.reclamation.retain",
|
||||
);
|
||||
}
|
||||
const executionClaim = execution
|
||||
? await withSessionEntryWorker(
|
||||
plan.databaseOptions,
|
||||
undefined,
|
||||
assertOpeningCurrent,
|
||||
(owner, source) =>
|
||||
owner.runExisting(source, async () => owner.captureGenerationClaim()),
|
||||
undefined,
|
||||
execution,
|
||||
signal,
|
||||
)
|
||||
: undefined;
|
||||
const claim = executionClaim ?? (retained?.found ? retained.claim : undefined);
|
||||
if (!claim) {
|
||||
throw new Error("Cold transcript operation lost its owning database");
|
||||
}
|
||||
const assertAllowed = () => {
|
||||
claim.assertCurrent();
|
||||
assertRequestCurrent();
|
||||
|
|
@ -117,7 +139,7 @@ async function runColdMutation(
|
|||
const diagnostics: SqliteSessionReclamationDiagnostics = { kind: plan.kind };
|
||||
const [completed] = await withSqliteReclamationAuthorization(
|
||||
commitGate,
|
||||
database.db,
|
||||
retained?.found ? retained.database.db : plan.databaseOptions.path,
|
||||
assertAllowed,
|
||||
(authorize) =>
|
||||
runSqliteTranscriptArchiveWorkerOperation<{
|
||||
|
|
@ -127,10 +149,12 @@ async function runColdMutation(
|
|||
diagnostics,
|
||||
signal,
|
||||
expectedMessageType: "reclaimed",
|
||||
validationOwner: { database, isCurrent: claim.isCurrent },
|
||||
onCommitRequest: () => {
|
||||
authorize();
|
||||
},
|
||||
validationOwner: executionClaim
|
||||
? { source: plan.databaseOptions, claim: executionClaim }
|
||||
: retained?.found
|
||||
? { database: retained.database, isCurrent: retained.claim.isCurrent }
|
||||
: undefined,
|
||||
onCommitRequest: authorize,
|
||||
withWriteAdmission: async (run, reclamationAdmission) =>
|
||||
runExclusiveSqliteSessionWrite(
|
||||
plan.databaseOptions,
|
||||
|
|
@ -173,17 +197,21 @@ async function runColdMutation(
|
|||
plan.kind === "cold-restore" &&
|
||||
completed.result.restored &&
|
||||
completed.result.sessionKey !== undefined &&
|
||||
claim.isCurrent()
|
||||
retained?.found &&
|
||||
retained.claim.isCurrent()
|
||||
) {
|
||||
assertAllowed();
|
||||
sessionChanges.emit({
|
||||
storePath: database.path,
|
||||
storePath: retained.database.path,
|
||||
sessionKey: completed.result.sessionKey,
|
||||
});
|
||||
}
|
||||
return completed.result;
|
||||
} finally {
|
||||
claim.release();
|
||||
await execution?.release();
|
||||
if (retained?.found) {
|
||||
retained.claim.release();
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
|
|
@ -205,134 +233,68 @@ type ColdBatchResult = SessionColdMaintenanceResult & {
|
|||
|
||||
async function archiveSessionColdBatch(options: ColdBatchOptions): Promise<ColdBatchResult> {
|
||||
const storePath = resolveOpenClawAgentSqlitePath(options.databaseOptions);
|
||||
const source = createOpenClawAgentDatabasePathMatcher();
|
||||
source(storePath, storePath);
|
||||
const assertCurrent = () => {
|
||||
options.assertCurrent?.();
|
||||
if (!source.isCurrent()) {
|
||||
throw new Error("Cold transcript database changed during maintenance");
|
||||
}
|
||||
};
|
||||
return operations.enqueue(storePath, async () => {
|
||||
const selection = await runExclusiveSqliteSessionWrite(
|
||||
options.databaseOptions,
|
||||
async () =>
|
||||
withOpenClawAgentDatabaseReadOnly((database) => {
|
||||
options.assertCurrent?.();
|
||||
const db = getNodeSqliteKysely<DB>(database.db);
|
||||
const excluded = readSessionColdStorageProtection(database, options.beforeMs);
|
||||
const admissions = collectAdmissionProtectedSessionIds({
|
||||
database,
|
||||
storePath: options.ownerStorePath,
|
||||
});
|
||||
for (const id of admissions) {
|
||||
excluded.add(id);
|
||||
}
|
||||
const cooled = new Set<string>();
|
||||
const now = Date.now();
|
||||
for (const cache of [restoredUntil, oversizedUntil]) {
|
||||
for (const [key, until] of cache) {
|
||||
if (until <= now) {
|
||||
cache.delete(key);
|
||||
} else if (key.startsWith(`${storePath}\0`)) {
|
||||
cooled.add(key.slice(storePath.length + 1));
|
||||
}
|
||||
}
|
||||
}
|
||||
for (const id of cooled) {
|
||||
excluded.add(id);
|
||||
}
|
||||
const externalizations = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("session_transcript_cold_archives")
|
||||
.select("session_id")
|
||||
.where("storage", "=", "sqlite")
|
||||
.$if(cooled.size + admissions.size > 0, (query) =>
|
||||
query.where("session_id", "not in", sqliteStringSet([...cooled, ...admissions])),
|
||||
)
|
||||
.orderBy("archived_at")
|
||||
.orderBy("session_id")
|
||||
.limit(options.maxTranscripts),
|
||||
).rows.flatMap((row) => {
|
||||
const archive = readSessionColdTranscript(database.db, row.session_id);
|
||||
return archive ? [archive] : [];
|
||||
});
|
||||
const candidates =
|
||||
externalizations.length < options.maxTranscripts
|
||||
? executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("session_windows as window")
|
||||
.leftJoin(
|
||||
"session_transcript_cold_archives as cold",
|
||||
"cold.session_id",
|
||||
"window.session_id",
|
||||
)
|
||||
.select("window.session_id")
|
||||
.where("cold.session_id", "is", null)
|
||||
.$if(excluded.size > 0, (query) =>
|
||||
query.where("window.session_id", "not in", sqliteStringSet([...excluded])),
|
||||
)
|
||||
.where("window.transcript_updated_at", "<", options.beforeMs)
|
||||
.where((eb) =>
|
||||
eb.exists(
|
||||
eb
|
||||
.selectFrom("transcript_events as event")
|
||||
.select("event.seq")
|
||||
.whereRef("event.session_id", "=", "window.session_id"),
|
||||
),
|
||||
)
|
||||
.orderBy("window.transcript_updated_at")
|
||||
.orderBy("window.session_id")
|
||||
.limit(options.maxTranscripts - externalizations.length),
|
||||
).rows
|
||||
: [];
|
||||
const databaseOptions = workerDatabaseOptions(options.databaseOptions);
|
||||
const plans = candidates.flatMap(({ session_id: sessionId }) => {
|
||||
const snapshot = readSessionStateDeleteSnapshot(database.db, sessionId);
|
||||
return snapshot.generation && snapshot.lastSeq !== null
|
||||
? [{ databaseOptions, sessionId, beforeMs: options.beforeMs, snapshot }]
|
||||
: [];
|
||||
});
|
||||
const freePages = Number(
|
||||
// sqlite-allow-raw -- Physical maintenance is needed only when SQLite owns free pages.
|
||||
database.db.prepare("PRAGMA freelist_count").get()?.freelist_count ?? 0,
|
||||
);
|
||||
return {
|
||||
input: {
|
||||
databaseOptions,
|
||||
plans,
|
||||
externalizations,
|
||||
maxBytes: options.maxBytes,
|
||||
} satisfies SessionColdBatchInput,
|
||||
freePages,
|
||||
};
|
||||
}, options.databaseOptions),
|
||||
"session.history.eviction-prepare",
|
||||
assertCurrent();
|
||||
const cooled = new Set<string>();
|
||||
const now = Date.now();
|
||||
for (const cache of [restoredUntil, oversizedUntil]) {
|
||||
for (const [key, until] of cache) {
|
||||
if (until <= now) {
|
||||
cache.delete(key);
|
||||
} else if (key.startsWith(`${storePath}\0`)) {
|
||||
cooled.add(key.slice(storePath.length + 1));
|
||||
}
|
||||
}
|
||||
}
|
||||
const input: SessionColdPreparationWorkerData["input"] = {
|
||||
databaseOptions: workerDatabaseOptions(options.databaseOptions),
|
||||
admissionIdentities: [
|
||||
...(collectActiveSessionWorkAdmissions().get(options.ownerStorePath) ?? []),
|
||||
],
|
||||
cooledSessionIds: [...cooled],
|
||||
beforeMs: options.beforeMs,
|
||||
maxTranscripts: options.maxTranscripts,
|
||||
maxBytes: options.maxBytes,
|
||||
};
|
||||
const [batch] = await withSqliteMutationWorkerLifetime(
|
||||
input.databaseOptions,
|
||||
async ({ assertCurrent: assertRequestCurrent, signal }) => {
|
||||
const prepared = await runSqliteTranscriptArchiveWorkerOperation<SessionColdBatchPrepared>({
|
||||
expectedMessageType: "done",
|
||||
signal,
|
||||
assertCurrent: () => {
|
||||
assertRequestCurrent();
|
||||
assertCurrent();
|
||||
},
|
||||
workerData: {
|
||||
type: "sqlite-transcript-archive-v2",
|
||||
operation: "cold-prepare",
|
||||
input,
|
||||
} satisfies SessionColdPreparationWorkerData,
|
||||
});
|
||||
assertRequestCurrent();
|
||||
assertCurrent();
|
||||
return prepared;
|
||||
},
|
||||
);
|
||||
assertCurrent();
|
||||
if (!batch) {
|
||||
throw new Error("Cold archive worker returned no prepared batch");
|
||||
}
|
||||
const empty: ColdBatchResult = {
|
||||
archivedTranscripts: 0,
|
||||
externalizedTranscripts: 0,
|
||||
envelopeBytes: 0,
|
||||
attemptedTranscripts: 0,
|
||||
};
|
||||
if (!selection.found) {
|
||||
return empty;
|
||||
}
|
||||
const { input, freePages } = selection.value;
|
||||
if (input.plans.length + input.externalizations.length === 0) {
|
||||
if (freePages > 0) {
|
||||
await runColdMutation(
|
||||
{ kind: "cold-maintain", databaseOptions: input.databaseOptions },
|
||||
options.assertCurrent,
|
||||
);
|
||||
}
|
||||
return empty;
|
||||
}
|
||||
const [batch] = await runSqliteTranscriptArchiveWorkerOperation<SessionColdBatchPrepared>({
|
||||
expectedMessageType: "done",
|
||||
workerData: {
|
||||
type: "sqlite-transcript-archive-v2",
|
||||
operation: "cold-prepare",
|
||||
input,
|
||||
} satisfies SessionColdPreparationWorkerData,
|
||||
});
|
||||
if (!batch) {
|
||||
throw new Error("Cold archive worker returned no prepared batch");
|
||||
}
|
||||
for (const sessionId of batch.oversizedSessionIds) {
|
||||
oversizedUntil.set(`${storePath}\0${sessionId}`, Date.now() + RESTORE_COOLDOWN_MS);
|
||||
log.warn("Transcript remains in SQLite because its archive exceeds the 64 MiB limit", {
|
||||
|
|
@ -352,24 +314,22 @@ async function archiveSessionColdBatch(options: ColdBatchOptions): Promise<ColdB
|
|||
prepared: batch.prepared,
|
||||
externalizations: batch.externalizations,
|
||||
beforeMs: options.beforeMs,
|
||||
protectionKeys: batch.protectionKeys,
|
||||
},
|
||||
() => {
|
||||
options.assertCurrent?.();
|
||||
const read = withOpenClawAgentDatabaseReadOnly(
|
||||
(database) =>
|
||||
collectAdmissionProtectedSessionIds({
|
||||
database,
|
||||
storePath: options.ownerStorePath,
|
||||
}),
|
||||
input.databaseOptions,
|
||||
);
|
||||
if (!read.found) {
|
||||
throw new Error("Cold transcript database disappeared");
|
||||
assertCurrent();
|
||||
const admissions = collectActiveSessionWorkAdmissions().get(options.ownerStorePath);
|
||||
if (
|
||||
[...(admissions ?? [])].some((identity) =>
|
||||
batch.protectionKeys.includes(normalizeStoreSessionKey(identity)),
|
||||
)
|
||||
) {
|
||||
throw new Error("Transcript became active; cold archival was canceled");
|
||||
}
|
||||
if (
|
||||
included.some(
|
||||
(id) =>
|
||||
read.value.has(id) ||
|
||||
admissions?.has(id) ||
|
||||
(restoredUntil.get(`${storePath}\0${id}`) ?? 0) > Date.now(),
|
||||
)
|
||||
) {
|
||||
|
|
@ -377,7 +337,13 @@ async function archiveSessionColdBatch(options: ColdBatchOptions): Promise<ColdB
|
|||
}
|
||||
},
|
||||
)
|
||||
: empty;
|
||||
: batch.freePages > 0
|
||||
? await runColdMutation(
|
||||
{ kind: "cold-maintain", databaseOptions: input.databaseOptions },
|
||||
assertCurrent,
|
||||
)
|
||||
: empty;
|
||||
assertCurrent();
|
||||
return {
|
||||
archivedTranscripts: result.archivedTranscripts,
|
||||
externalizedTranscripts: result.externalizedTranscripts,
|
||||
|
|
@ -507,18 +473,63 @@ export async function restoreSessionColdTranscript(
|
|||
});
|
||||
}
|
||||
|
||||
function configuredStores(
|
||||
config: OpenClawConfig,
|
||||
): Array<{ agentId: string; storePath: string; ownerStorePath: string }> {
|
||||
return resolveSessionStoreTargets(config, { allAgents: true }).map((target) => {
|
||||
const resolved = resolveSqliteTranscriptReadScope({ ...target, sessionId: "cold-maintenance" });
|
||||
const options = toDatabaseOptions(resolved);
|
||||
return {
|
||||
agentId: options.agentId,
|
||||
storePath: resolveOpenClawAgentSqlitePath(options),
|
||||
ownerStorePath: target.storePath,
|
||||
};
|
||||
async function configuredStores(config: OpenClawConfig, assertCallerCurrent?: () => void) {
|
||||
const prepared = prepareSessionStoreTargetInventory(
|
||||
config,
|
||||
listConfiguredSessionStoreAgentIds(config),
|
||||
process.env,
|
||||
"configured",
|
||||
);
|
||||
const { env, candidates } = prepared;
|
||||
const context = captureOpenClawStateReadWorkerContext({ env });
|
||||
const registryRead = prepareOpenClawAgentDatabaseRegistrySnapshotRead({ env });
|
||||
const source = createOpenClawAgentDatabasePathMatcher();
|
||||
for (const candidate of candidates) {
|
||||
source(candidate.path, candidate.path);
|
||||
}
|
||||
const assertCurrent = () => {
|
||||
assertCallerCurrent?.();
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
context.admission.assertCurrent();
|
||||
if (!source.isCurrent()) {
|
||||
throw new Error("Session store changed during cold maintenance");
|
||||
}
|
||||
};
|
||||
const stores = await withSessionHistoryWorkerReadCandidates(candidates, async (discovery) => {
|
||||
const registry = await registryRead.read();
|
||||
assertCurrent();
|
||||
registryRead.assertCurrent();
|
||||
discovery.assertCurrent();
|
||||
const inventory = await discovery.readTargetInventory({
|
||||
...prepared,
|
||||
registeredDatabases:
|
||||
registry.result.status === "available"
|
||||
? registry.result.entries
|
||||
: { status: "unavailable" },
|
||||
});
|
||||
assertCurrent();
|
||||
registryRead.assertCurrent();
|
||||
discovery.assertCurrent();
|
||||
if (inventory.kind === "session-target-registry-required") {
|
||||
throw new Error("Cold maintenance did not receive its registry snapshot");
|
||||
}
|
||||
return inventory.agents.flatMap(({ result, reads }) =>
|
||||
result.available
|
||||
? result.targets.map((target, index) => ({
|
||||
databaseOptions: {
|
||||
...reads[index]!.database,
|
||||
// Archive paths and restore cooldowns retain the configured SQLite locator.
|
||||
path: reads[index]!.target.storePath,
|
||||
env,
|
||||
},
|
||||
ownerStorePath: target.storePath,
|
||||
}))
|
||||
: [],
|
||||
);
|
||||
});
|
||||
assertCurrent();
|
||||
registryRead.assertCurrent();
|
||||
return { stores, assertCurrent };
|
||||
}
|
||||
|
||||
export async function runSessionColdStorageMaintenance(params: {
|
||||
|
|
@ -535,31 +546,29 @@ export async function runSessionColdStorageMaintenance(params: {
|
|||
return result;
|
||||
}
|
||||
const beforeMs = Date.now() - (config.afterDays ?? 30) * 24 * 60 * 60 * 1000;
|
||||
const stores = configuredStores(params.config);
|
||||
const { stores, assertCurrent } = await configuredStores(params.config, params.assertCurrent);
|
||||
assertCurrent();
|
||||
const start = stores.length ? nextStore % stores.length : 0;
|
||||
nextStore = start + 1;
|
||||
let remainingTranscripts = MAX_TRANSCRIPTS_PER_PASS;
|
||||
let remainingBytes = MAX_BATCH_BYTES;
|
||||
for (const { agentId, storePath, ownerStorePath } of [
|
||||
for (const { databaseOptions, ownerStorePath } of [
|
||||
...stores.slice(start),
|
||||
...stores.slice(0, start),
|
||||
]) {
|
||||
if (remainingTranscripts <= 0 || remainingBytes <= 0) {
|
||||
break;
|
||||
}
|
||||
params.assertCurrent?.();
|
||||
const exists = withOpenClawAgentDatabaseReadOnly(() => true, { agentId, path: storePath });
|
||||
if (!exists.found) {
|
||||
continue;
|
||||
}
|
||||
assertCurrent();
|
||||
const batch = await archiveSessionColdBatch({
|
||||
databaseOptions: { agentId, path: storePath },
|
||||
databaseOptions,
|
||||
ownerStorePath,
|
||||
beforeMs,
|
||||
maxTranscripts: remainingTranscripts,
|
||||
maxBytes: remainingBytes,
|
||||
assertCurrent: params.assertCurrent,
|
||||
assertCurrent,
|
||||
});
|
||||
assertCurrent();
|
||||
result.archivedTranscripts += batch.archivedTranscripts;
|
||||
result.externalizedTranscripts += batch.externalizedTranscripts;
|
||||
remainingTranscripts -= batch.attemptedTranscripts;
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue