test(ci): trim duplicate SQLite lifecycle tours (#122308)

* test(ci): accelerate SQLite flip proof

* test(sessions): replace duplicate archive stress tour

* test(sessions): preserve SQLite lifecycle contracts

* test(sessions): drop duplicate archive assertions

---------

Co-authored-by: Amp <amp@ampcode.com>
This commit is contained in:
Peter Steinberger
2026-08-11 16:51:28 -07:00
committed by GitHub
parent 087fb56f77
commit 76c7cd7308
4 changed files with 169 additions and 839 deletions
@@ -67,36 +67,6 @@ export function assertSqliteFlipProofCore(report: SqliteFlipProofReport): void {
),
),
).toBe(true);
expect(report.pluginSdkConsumer).toMatchObject({
activeJsonlForSessionExists: false,
latestAssistantTextBeforeAppend: report.fullTurnAssistantText,
latestAssistantTextAfterAppend: "sqlite sdk consumer appended by identity",
sessionKey: report.pluginSdkSessionKey,
});
expect(report.pluginSdkConsumer?.sessionIdentity).toBe(report.pluginSdkSessionKey);
expect(report.pluginSdkConsumer?.listedSessionKeys).toContain(report.pluginSdkSessionKey);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-plugin-sdk-consumer" &&
checkpoint.sqlite.trackedEntries.some(
(entry) => entry.sessionKey === report.pluginSdkSessionKey && entry.transcriptEvents >= 3,
),
),
).toBe(true);
const cleanupCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-cleanup-pruning",
);
expect(
cleanupCheckpoint?.sqlite.trackedEntries.some(
(entry) => entry.sessionKey === report.cleanupPruneSessionKey,
),
).toBe(false);
const cleanupArchive = cleanupCheckpoint?.archiveArtifacts.find(
(artifact) =>
artifact.archiveReason === "deleted" && artifact.archiveSessionId === "sqlite-cleanup-prune",
);
expect(cleanupArchive?.messageTexts).toContain("sqlite cleanup prune me");
const idempotenceCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-doctor-import-idempotence",
);
@@ -145,4 +115,164 @@ export function assertSqliteFlipProofCore(report: SqliteFlipProofReport): void {
(entry) => entry.sessionKey === report.concurrentDeleteSessionKey,
),
).toBe(false);
expect(report.checkpoints.map((checkpoint) => checkpoint.label)).toEqual([
"seeded-legacy-store",
"after-startup-import",
"after-doctor-inspect",
"after-doctor-validate",
"after-rollback-restore",
"after-gateway-restart",
"after-chat-send",
"after-full-agent-turn",
"after-doctor-import-idempotence",
"after-downgrade-reupgrade-import",
"after-sqlite-busy-contention",
"after-concurrent-multi-client",
"after-sessions-reset",
"after-second-startup-after-reset",
"after-transcript-append",
"after-sessions-delete",
"after-shared-first-delete",
"after-shared-final-delete",
"after-final-doctor-inspect",
]);
expect(
startupImportCheckpoint?.archiveArtifacts.some(
(artifact) =>
artifact.path.includes("old-orphan.deleted.jsonl") &&
artifact.textTail?.includes("old-orphan") === true,
),
).toBe(true);
expect(report.rollbackRestore).toMatchObject({
archivedBeforeRestore: true,
failedManifestIssueCode: "e2e_forced_post_archive_failure",
sourceRestored: true,
sqliteStillExists: true,
});
expect(report.rollbackRestore?.manifestPath).toContain("session-sqlite-migration-runs");
expect(
report.rollbackRestore?.restoredFiles.some((filePath) =>
filePath.replaceAll("\\", "/").endsWith("/sqlite-rollback-restore.jsonl"),
),
).toBe(true);
expect(
report.rollbackRestore?.idempotentRestoreSkippedFiles.some((filePath) =>
filePath.replaceAll("\\", "/").endsWith("/sqlite-rollback-restore.jsonl"),
),
).toBe(true);
expect(report.scaleMigration).toMatchObject({
minTranscriptEventsPerSession: 4,
seededEvents: 96,
seededSessions: 24,
});
expect(report.scaleMigration?.importedSessionKeys).toHaveLength(24);
expect(report.scaleMigration?.startupImportElapsedMs).toBeGreaterThanOrEqual(0);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-full-agent-turn" &&
checkpoint.sqlite.trackedEntries.some(
(entry) =>
entry.sessionKey === report.fullTurnSessionKey &&
entry.transcriptEvents >= 2 &&
entry.trajectoryEvents >= 1,
),
),
).toBe(true);
expect(report.downgradeReupgrade).toMatchObject({
activeJsonlArchived: true,
doctorImportedEntries: 1,
doctorImportedTranscriptEvents: 2,
sessionId: "sqlite-downgrade-reupgrade",
sessionKey: "agent:main:dashboard:sqlite-downgrade-reupgrade",
trajectoryPointerArchived: true,
trajectoryPointerSourceRemoved: true,
trajectorySidecarArchived: true,
trajectorySidecarSourceRemoved: true,
transcriptEvents: 2,
});
const downgradeCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-downgrade-reupgrade-import",
);
expect(
downgradeCheckpoint?.archiveArtifacts.some(
(artifact) =>
artifact.path.includes("sqlite-downgrade-reupgrade.trajectory.jsonl") &&
artifact.textTail?.includes("trajectory") === true,
),
).toBe(true);
expect(
downgradeCheckpoint?.archiveArtifacts.some((artifact) =>
artifact.path.includes("sqlite-downgrade-reupgrade.trajectory-path.json"),
),
).toBe(true);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-downgrade-reupgrade-import" &&
checkpoint.sqlite.trackedEntries.some(
(entry) =>
entry.sessionKey === "agent:main:dashboard:sqlite-downgrade-reupgrade" &&
entry.transcriptEvents === 2,
),
),
).toBe(true);
expect(report.busyContention).toMatchObject({
childExitCode: 0,
childSignal: null,
holdMs: 500,
sessionId: "sqlite-busy-contention",
sessionKey: "agent:main:dashboard:sqlite-busy-contention",
transcriptEvents: 2,
});
expect(report.busyContention?.elapsedMs).toBeGreaterThanOrEqual(250);
expect(report.secondStartupAfterReset).toMatchObject({
activeJsonlForSessionExists: false,
historyContainsPostResetAppend: true,
sessionKey: report.resetSessionKey,
});
expect(report.secondStartupAfterReset?.transcriptEvents).toBeGreaterThanOrEqual(1);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-transcript-append" &&
checkpoint.sqlite.trackedEntries.some(
(entry) => entry.sessionKey === report.resetSessionKey && entry.transcriptEvents >= 1,
),
),
).toBe(true);
const deleteCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-sessions-delete",
);
const deleteArchive = deleteCheckpoint?.archiveArtifacts.find(
(artifact) =>
artifact.archiveReason === "deleted" && artifact.archiveSessionId === "sqlite-delete-session",
);
expect(deleteArchive?.messageTexts).toContain("delete me");
const sharedFinalCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-shared-final-delete",
);
const sharedFinalArchive = sharedFinalCheckpoint?.archiveArtifacts.find(
(artifact) =>
artifact.archiveReason === "deleted" && artifact.archiveSessionId === "sqlite-shared-session",
);
const retainedSharedImportSources = sharedFinalCheckpoint?.archiveArtifacts.filter(
(artifact) =>
artifact.path.includes("session-sqlite-import-archive") &&
(artifact.path.includes("sqlite-shared-a.jsonl") ||
artifact.path.includes("sqlite-shared-b.jsonl")),
);
expect(
sharedFinalArchive?.messageTexts?.includes("shared") ||
(retainedSharedImportSources?.length === 2 &&
retainedSharedImportSources.every((artifact) =>
artifact.messageTexts?.some((text) => text.includes("shared")),
)),
).toBe(true);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-shared-final-delete" && checkpoint.archiveArtifacts.length > 0,
),
).toBe(true);
}
@@ -25,21 +25,9 @@ import {
connectGatewayClient,
disconnectGatewayClient,
} from "../../src/gateway/test-helpers.e2e.js";
import {
getSessionEntry as getSdkSessionEntry,
listSessionEntries as listSdkSessionEntries,
loadTranscriptEventsSync as loadSdkTranscriptEventsSync,
} from "../../src/plugin-sdk/session-store-runtime.js";
import {
appendSessionTranscriptMessageByIdentity,
readLatestAssistantTextByIdentity,
readSessionTranscriptEvents,
resolveSessionTranscriptIdentity,
} from "../../src/plugin-sdk/session-transcript-runtime.js";
import { closeOpenClawAgentDatabasesForTest } from "../../src/state/openclaw-agent-db.js";
import { closeOpenClawStateDatabaseForTest } from "../../src/state/openclaw-state-db.js";
import { sleep } from "../../src/utils.js";
import { normalizeSessionDeliveryState } from "../../src/utils/delivery-context.shared.js";
import { createOpenClawTestInstance } from "./openclaw-test-instance.js";
type DoctorMode = "import" | "inspect" | "validate" | "restore";
@@ -52,8 +40,6 @@ type DoctorCommandEvidence = Awaited<ReturnType<typeof runDoctor>>;
type FileInventoryEntry = Awaited<ReturnType<typeof inventoryFile>>;
type ProofCheckpoint = Awaited<ReturnType<typeof captureCheckpoint>>;
type PluginSdkConsumerEvidence = Awaited<ReturnType<typeof runPluginSdkConsumerProbe>>;
type ManualCompactionEvidence = Awaited<ReturnType<typeof runManualCompactionProof>>;
type ScaleMigrationEvidence = ReturnType<typeof requireScaleMigrationProof>;
type DowngradeReupgradeEvidence = Awaited<ReturnType<typeof runDowngradeReupgradeProof>>;
type BusyContentionEvidence = Awaited<ReturnType<typeof runSqliteBusyContentionProof>>;
@@ -77,14 +63,8 @@ const CONCURRENT_RESET_SESSION_KEY = "agent:main:dashboard:sqlite-concurrent-res
const CONCURRENT_DELETE_SESSION_KEY = "agent:main:dashboard:sqlite-concurrent-delete";
const CONCURRENT_SEND_TEXT = "sqlite concurrent send history reset";
const CONCURRENT_DELETE_TEXT = "sqlite concurrent delete while send is active";
const CLEANUP_PRUNE_SESSION_ID = "sqlite-cleanup-prune";
const CLEANUP_PRUNE_SESSION_KEY = "agent:main:dashboard:sqlite-cleanup-prune";
const CLEANUP_PRUNE_TEXT = "sqlite cleanup prune me";
const FULL_TURN_ASSISTANT_TEXT = "OPENCLAW_E2E_OK_12";
const FULL_TURN_SESSION_KEY = "agent:main:sqlite-full-turn";
const MANUAL_COMPACTION_SESSION_KEY = "agent:main:dashboard:sqlite-manual-compact";
const PLUGIN_SDK_APPEND_TEXT = "sqlite sdk consumer appended by identity";
const PLUGIN_SDK_SESSION_KEY = "agent:main:dashboard:sqlite-sdk-consumer";
const DOWNGRADE_REUPGRADE_SESSION_ID = "sqlite-downgrade-reupgrade";
const DOWNGRADE_REUPGRADE_SESSION_KEY = "agent:main:dashboard:sqlite-downgrade-reupgrade";
const DOWNGRADE_REUPGRADE_TEXT = "sqlite downgrade wrote file-backed state";
@@ -120,6 +100,13 @@ export async function runSqliteSessionsTranscriptsFlipProof(options: RunOptions
HTTP_PROXY: undefined,
HTTPS_PROXY: undefined,
NO_PROXY: "127.0.0.1,localhost",
...(options.requireBuiltCli !== true
? {
OPENCODE_API_KEY: undefined,
OPENCODE_ZEN_API_KEY: undefined,
OPENCLAW_DISABLE_BUNDLED_PLUGINS: "1",
}
: {}),
OPENAI_API_KEY: "sk-openclaw-e2e-mock",
OPENCLAW_TEST_MINIMAL_GATEWAY: undefined,
OPENCLAW_SKIP_PROVIDERS: undefined,
@@ -128,8 +115,7 @@ export async function runSqliteSessionsTranscriptsFlipProof(options: RunOptions
startTimeoutMs: 90_000,
stopTimeoutMs: 3_000,
});
// The proof mixes child-process CLI calls with in-process SDK accessors.
// Apply the fixture environment so both paths resolve the same isolated DB.
// Doctor commands and the gateway must resolve the same isolated database.
inst.state.applyEnv();
const context = buildProofContext(inst.stateDir);
const checkpoints: ProofCheckpoint[] = [];
@@ -137,9 +123,7 @@ export async function runSqliteSessionsTranscriptsFlipProof(options: RunOptions
let gatewayEntrypoint: string[] = [];
let busyContention: BusyContentionEvidence | undefined;
let downgradeReupgrade: DowngradeReupgradeEvidence | undefined;
let manualCompaction: ManualCompactionEvidence | undefined;
let mockOpenAi: ProofChildProcess | undefined;
let pluginSdkConsumer: PluginSdkConsumerEvidence | undefined;
let rollbackRestore: RollbackRestoreEvidence | undefined;
let scaleMigration: ScaleMigrationEvidence | undefined;
let secondStartupAfterReset: SecondStartupAfterResetEvidence | undefined;
@@ -249,37 +233,6 @@ export async function runSqliteSessionsTranscriptsFlipProof(options: RunOptions
await requireMockOpenAiRequest(context.mockOpenAiRequestLog);
await record("after-full-agent-turn");
manualCompaction = await runManualCompactionProof(restartedClient, context);
await record("after-manual-compaction");
const pluginSdkRunId = await sendGatewayUserMessage(
restartedClient,
context.pluginSdkSessionKey,
`Reply with exactly ${context.fullTurnAssistantText}. SDK consumer proof.`,
);
await waitForAgentRunOk(restartedClient, pluginSdkRunId);
const pluginSdkSessionId = await waitForSqliteSessionId(
context.agentDbPath,
context.pluginSdkSessionKey,
);
await waitForSqliteMessageContains(
context.agentDbPath,
pluginSdkSessionId,
"assistant",
context.fullTurnAssistantText,
);
pluginSdkConsumer = await runPluginSdkConsumerProbe(context, pluginSdkSessionId);
await waitForSqliteMessageContains(
context.agentDbPath,
pluginSdkSessionId,
"assistant",
context.pluginSdkAppendText,
);
await record("after-plugin-sdk-consumer");
await runGatewayCleanupPruningProof(restartedClient, context);
await record("after-cleanup-pruning");
const idempotentImportDoctor = await runDoctorIdempotenceProof(inst, context);
await record("after-doctor-import-idempotence", idempotentImportDoctor);
@@ -357,7 +310,6 @@ export async function runSqliteSessionsTranscriptsFlipProof(options: RunOptions
ok: failures.length === 0,
agentId: context.agentId,
checkpoints,
cleanupPruneSessionKey: context.cleanupPruneSessionKey,
concurrentDeleteSessionKey: context.concurrentDeleteSessionKey,
concurrentResetSessionKey: context.concurrentResetSessionKey,
concurrentSendSessionKey: context.concurrentSendSessionKey,
@@ -367,12 +319,8 @@ export async function runSqliteSessionsTranscriptsFlipProof(options: RunOptions
fullTurnSessionKey: context.fullTurnSessionKey,
gatewayEntrypoint,
legacySessionId: context.legacySessionId,
...(manualCompaction ? { manualCompaction } : {}),
manualCompactionSessionKey: context.manualCompactionSessionKey,
mockOpenAiRequestLog: context.mockOpenAiRequestLog,
oldStateSessionKeys: [...context.oldStateSessionKeys],
...(pluginSdkConsumer ? { pluginSdkConsumer } : {}),
pluginSdkSessionKey: context.pluginSdkSessionKey,
resetSessionKey: context.resetSessionKey,
...(rollbackRestore ? { rollbackRestore } : {}),
...(busyContention ? { busyContention } : {}),
@@ -402,7 +350,6 @@ function buildProofContext(stateDir: string) {
agentDbPath: path.join(agentDir, "agent", "openclaw-agent.sqlite"),
agentId: AGENT_ID,
archiveRoots: [path.join(agentDir, "session-sqlite-import-archive"), activeSessionsDir],
cleanupPruneSessionKey: CLEANUP_PRUNE_SESSION_KEY,
concurrentDeleteSessionKey: CONCURRENT_DELETE_SESSION_KEY,
concurrentResetSessionKey: CONCURRENT_RESET_SESSION_KEY,
concurrentSendSessionKey: CONCURRENT_SEND_SESSION_KEY,
@@ -411,11 +358,8 @@ function buildProofContext(stateDir: string) {
fullTurnSessionKey: FULL_TURN_SESSION_KEY,
legacySessionsDir,
legacySessionId: "sqlite-legacy-main",
manualCompactionSessionKey: MANUAL_COMPACTION_SESSION_KEY,
mockOpenAiRequestLog: path.join(stateDir, "mock-openai-requests.ndjson"),
oldStateSessionKeys: [...OLD_STATE_SESSION_KEYS],
pluginSdkAppendText: PLUGIN_SDK_APPEND_TEXT,
pluginSdkSessionKey: PLUGIN_SDK_SESSION_KEY,
resetSessionKey: RESET_SESSION_KEY,
sharedSessionKeys: [...SHARED_SESSION_KEYS],
stateDir,
@@ -423,13 +367,10 @@ function buildProofContext(stateDir: string) {
trackedSessionKeys: [
RESET_SESSION_KEY,
DELETE_SESSION_KEY,
CLEANUP_PRUNE_SESSION_KEY,
CONCURRENT_SEND_SESSION_KEY,
CONCURRENT_RESET_SESSION_KEY,
CONCURRENT_DELETE_SESSION_KEY,
FULL_TURN_SESSION_KEY,
MANUAL_COMPACTION_SESSION_KEY,
PLUGIN_SDK_SESSION_KEY,
DOWNGRADE_REUPGRADE_SESSION_KEY,
SQLITE_BUSY_SESSION_KEY,
...SHARED_SESSION_KEYS,
@@ -991,192 +932,6 @@ async function appendProofMessage(
}
}
async function runManualCompactionProof(client: GatewayClient, context: ProofContext) {
const runId = await sendGatewayUserMessage(
client,
context.manualCompactionSessionKey,
`Reply with exactly ${context.fullTurnAssistantText}. Manual compaction proof.`,
);
await waitForAgentRunOk(client, runId);
const sessionId = await waitForSqliteSessionId(
context.agentDbPath,
context.manualCompactionSessionKey,
);
await waitForSqliteMessageContains(
context.agentDbPath,
sessionId,
"assistant",
context.fullTurnAssistantText,
);
const rowCountBefore = countSqliteTranscriptEvents(context.agentDbPath, sessionId);
if (rowCountBefore < 2) {
throw new Error(
`manual compaction source transcript had too few rows: ${rowCountBefore} for ${sessionId}`,
);
}
const listed: { sessions?: Array<{ key?: string; sessionId?: string }> } = await client.request(
"sessions.list",
{},
);
const listedSession = (listed.sessions ?? []).find(
(session) =>
session.sessionId === sessionId || session.key === context.manualCompactionSessionKey,
);
if (!listedSession?.key) {
throw new Error(
`manual compaction session was not listed before compact: ${JSON.stringify(listed)}`,
);
}
const compacted: {
compacted?: boolean;
key?: string;
ok?: boolean;
} = await client.request("sessions.compact", {
key: listedSession.key,
});
if (compacted.ok !== true || compacted.compacted !== true) {
throw new Error(
`manual compaction did not compact using ${listedSession.key}: ${JSON.stringify(
compacted,
)}; listed=${JSON.stringify(listed)}`,
);
}
const evidence = readSqliteEvidence(context.agentDbPath, [context.manualCompactionSessionKey]);
const row = evidence.trackedEntries.find(
(entry) => entry.sessionKey === context.manualCompactionSessionKey,
);
if (!row?.entry) {
throw new Error(`manual compaction entry missing for ${context.manualCompactionSessionKey}`);
}
const checkpointCount = Array.isArray(row.entry.compactionCheckpoints)
? row.entry.compactionCheckpoints.length
: 0;
if (checkpointCount < 1) {
throw new Error(`manual compaction did not write checkpoint metadata: ${JSON.stringify(row)}`);
}
if (Object.hasOwn(row.entry, "sessionFile")) {
throw new Error(`manual compaction entry retained file-era identity: ${JSON.stringify(row)}`);
}
return {
checkpointCount,
compacted: compacted.compacted,
rowCountAfter: countSqliteTranscriptEvents(context.agentDbPath, row.sessionId),
rowCountBefore,
transcriptIdentity: context.manualCompactionSessionKey,
sessionId: row.sessionId,
sessionKey: context.manualCompactionSessionKey,
};
}
async function runPluginSdkConsumerProbe(context: ProofContext, sessionId: string) {
const scope = {
agentId: context.agentId,
sessionId,
sessionKey: context.pluginSdkSessionKey,
storePath: context.storePath,
};
const sessionEntry = getSdkSessionEntry({
agentId: context.agentId,
readConsistency: "latest",
sessionKey: context.pluginSdkSessionKey,
storePath: context.storePath,
});
if (sessionEntry?.sessionId !== sessionId) {
throw new Error(
`SDK session store read returned ${JSON.stringify(sessionEntry)} for ${context.pluginSdkSessionKey}`,
);
}
if (Object.hasOwn(sessionEntry, "sessionFile")) {
throw new Error(`SDK session store exposed retired transcript locator`);
}
const listedSessionKeys = listSdkSessionEntries({
agentId: context.agentId,
storePath: context.storePath,
}).map((entry) => entry.sessionKey);
if (!listedSessionKeys.includes(context.pluginSdkSessionKey)) {
throw new Error(`SDK session list omitted ${context.pluginSdkSessionKey}`);
}
const identity = await resolveSessionTranscriptIdentity(scope);
const latestBefore = await readLatestAssistantTextByIdentity(scope);
if (latestBefore?.text !== context.fullTurnAssistantText) {
throw new Error(
`SDK latest assistant read returned ${JSON.stringify(latestBefore)} for ${context.pluginSdkSessionKey}`,
);
}
const transcriptEventsBeforeAppend = (await readSessionTranscriptEvents(scope)).length;
const storeTranscriptEvents = loadSdkTranscriptEventsSync(scope).length;
const artifacts = sessionArtifactPaths(context.activeSessionsDir, sessionId);
const activeJsonlForSessionExists = fsSync.existsSync(artifacts.jsonl);
if (activeJsonlForSessionExists) {
throw new Error(`SDK probe found active JSONL for SQLite session at ${artifacts.jsonl}`);
}
const activeTrajectorySessionSidecarForSessionExists = fsSync.existsSync(artifacts.trajectory);
const activeTrajectoryPointerForSessionExists = fsSync.existsSync(artifacts.pointer);
const activeTrajectoryRuntimeSidecarForSessionExists = fsSync.existsSync(artifacts.runtime);
if (
activeTrajectorySessionSidecarForSessionExists ||
activeTrajectoryPointerForSessionExists ||
activeTrajectoryRuntimeSidecarForSessionExists
) {
throw new Error(
`SDK trajectory probe found active sidecar paths: ${JSON.stringify({
pointer: artifacts.pointer,
runtime: artifacts.runtime,
session: artifacts.trajectory,
})}`,
);
}
const appended = await appendSessionTranscriptMessageByIdentity({
...scope,
message: {
role: "assistant",
content: [{ type: "text", text: context.pluginSdkAppendText }],
timestamp: Date.now(),
},
});
if (!appended?.appended || !appended.messageId) {
throw new Error(`SDK transcript append failed for ${context.pluginSdkSessionKey}`);
}
const latestAfter = await readLatestAssistantTextByIdentity(scope);
if (latestAfter?.text !== context.pluginSdkAppendText) {
throw new Error(
`SDK latest assistant after append returned ${JSON.stringify(latestAfter)} for ${
context.pluginSdkSessionKey
}`,
);
}
const transcriptEventsAfterAppend = (await readSessionTranscriptEvents(scope)).length;
if (transcriptEventsAfterAppend <= transcriptEventsBeforeAppend) {
throw new Error(
`SDK transcript append did not increase event count for ${context.pluginSdkSessionKey}`,
);
}
return {
activeJsonlForSessionExists,
activeTrajectoryPointerForSessionExists,
activeTrajectoryRuntimeSidecarForSessionExists,
activeTrajectorySessionSidecarForSessionExists,
appendedMessageId: appended.messageId,
identityMemoryKey: identity.memoryKey,
latestAssistantTextBeforeAppend: latestBefore.text,
latestAssistantTextAfterAppend: latestAfter.text,
listedSessionKeys,
sessionIdentity: context.pluginSdkSessionKey,
sessionId,
sessionKey: context.pluginSdkSessionKey,
storeTranscriptEvents,
transcriptEventsAfterAppend,
transcriptEventsBeforeAppend,
};
}
function sessionArtifactPaths(sessionsDir: string, sessionId: string) {
return {
jsonl: path.join(sessionsDir, `${sessionId}.jsonl`),
@@ -1186,48 +941,6 @@ function sessionArtifactPaths(sessionsDir: string, sessionId: string) {
};
}
async function runGatewayCleanupPruningProof(
client: GatewayClient,
context: ProofContext,
): Promise<void> {
const updatedAt = Date.now() - 31 * 24 * 60 * 60 * 1000;
await importProofSession(
context,
context.cleanupPruneSessionKey,
CLEANUP_PRUNE_SESSION_ID,
{
delivery: normalizeSessionDeliveryState({ context: { channel: "cli" } }),
chatType: "direct",
sessionFile: formatSqliteSessionFileMarker({
agentId: context.agentId,
sessionId: CLEANUP_PRUNE_SESSION_ID,
storePath: context.storePath,
}),
sessionId: CLEANUP_PRUNE_SESSION_ID,
sessionStartedAt: updatedAt - 500,
updatedAt,
},
[messageEvent("sqlite-cleanup-prune-1", "user", CLEANUP_PRUNE_TEXT)],
);
await waitForSqliteMessageContains(
context.agentDbPath,
CLEANUP_PRUNE_SESSION_ID,
"user",
CLEANUP_PRUNE_TEXT,
);
const result: { afterCount?: number; applied?: boolean; pruned?: number } = await client.request(
"sessions.cleanup",
{ enforce: true },
{ timeoutMs: SQLITE_FLIP_PROOF_OPERATION_TIMEOUT_MS },
);
if (result?.applied !== true || (result.pruned ?? 0) < 1) {
throw new Error(`sessions.cleanup did not prune stale SQLite rows: ${JSON.stringify(result)}`);
}
await waitForSessionEntryAbsent(context.agentDbPath, context.cleanupPruneSessionKey);
await waitForSqliteEventsAbsent(context.agentDbPath, CLEANUP_PRUNE_SESSION_ID);
}
async function runDoctorIdempotenceProof(
inst: OpenClawTestInstance,
context: ProofContext,
@@ -1792,15 +1505,6 @@ async function waitForSessionEntryAbsent(dbPath: string, sessionKey: string): Pr
);
}
async function waitForSqliteEventsAbsent(dbPath: string, sessionId: string): Promise<void> {
await pollUntil(
() => countSqliteTranscriptEvents(dbPath, sessionId),
(eventCount) => eventCount === 0,
() => new Error(`timed out waiting for SQLite transcript row deletion for ${sessionId}`),
{ stableMs: 1_500 },
);
}
async function waitForSqliteMessageContains(
dbPath: string,
sessionId: string,
@@ -2229,21 +1933,6 @@ function validateCheckpointInvariants(
sessionId: "sqlite-delete-session",
});
}
if (checkpoint.label === "after-cleanup-pruning") {
requireArchiveText(checkpoint, failures, {
description: "cleanup-pruned transcript archive",
includes: [CLEANUP_PRUNE_TEXT],
reason: "deleted",
sessionId: CLEANUP_PRUNE_SESSION_ID,
});
if (
checkpoint.sqlite.trackedEntries.some(
(entry) => entry.sessionKey === context.cleanupPruneSessionKey,
)
) {
failures.push(`${checkpoint.label}: cleanup-pruned entry still exists in SQLite`);
}
}
if (checkpoint.label === "after-doctor-import-idempotence") {
const totals = checkpoint.doctor?.totals ?? {};
if (totals.importedEntries !== 0 || totals.importedTranscriptEvents !== 0) {
@@ -1,25 +1,5 @@
// Built-CLI SQLite flip proof requires dist entrypoints before running the gateway lifecycle.
import { createHash, randomBytes, randomUUID } from "node:crypto";
import fs from "node:fs";
import path from "node:path";
import { performance } from "node:perf_hooks";
import { DatabaseSync } from "node:sqlite";
import { describe, expect, it } from "vitest";
import { readSessionArchiveContentSync } from "../../src/config/sessions/archive-compression.js";
import {
loadSessionEntry,
loadTranscriptEvents,
replaceSessionEntry,
} from "../../src/config/sessions/session-accessor.js";
import { replaceTranscriptEvents } from "../../src/config/sessions/session-accessor.sqlite-transcript-write.js";
import { resolveSqliteTargetFromSessionStorePath } from "../../src/config/sessions/session-sqlite-target.js";
import {
connectGatewayClient,
disconnectGatewayClient,
} from "../../src/gateway/test-helpers.e2e.js";
import { closeOpenClawAgentDatabasesForTest } from "../../src/state/openclaw-agent-db.js";
import { closeOpenClawStateDatabaseForTest } from "../../src/state/openclaw-state-db.js";
import { createOpenClawTestInstance } from "../helpers/openclaw-test-instance.js";
import { assertSqliteFlipProofCore } from "../helpers/sqlite-sessions-transcripts-flip-proof-assertions.ts";
import { runSqliteSessionsTranscriptsFlipProof } from "../helpers/sqlite-sessions-transcripts-flip-proof.ts";
@@ -32,251 +12,4 @@ describe("SQLite sessions/transcripts flip built CLI proof", () => {
);
assertSqliteFlipProofCore(report);
}, 420_000);
it("keeps built gateway RPC responsive while deleting a large transcript", async () => {
const inst = await createOpenClawTestInstance({
name: `sqlite-archive-responsive-${randomUUID()}`,
startTimeoutMs: 90_000,
stopTimeoutMs: 5_000,
});
inst.state.applyEnv();
const sessionId = "sqlite-large-archive-responsive";
const sessionKey = "agent:main:dashboard:sqlite-large-archive-responsive";
const writerSessionKey = "agent:main:dashboard:sqlite-large-archive-writer";
const warmupSessionId = "sqlite-archive-worker-warmup";
const warmupSessionKey = "agent:main:dashboard:sqlite-archive-worker-warmup";
const storePath = path.join(inst.stateDir, "agents", "main", "sessions", "sessions.json");
const archiveDirectory = path.dirname(storePath);
const events = createLargeTranscriptEvents(sessionId);
const expectedArchiveContent = `${events.map((event) => JSON.stringify(event)).join("\n")}\n`;
let deleteClient: Awaited<ReturnType<typeof connectGatewayClient>> | undefined;
let probeClient: Awaited<ReturnType<typeof connectGatewayClient>> | undefined;
try {
await replaceSessionEntry({ sessionKey, storePath }, { sessionId, updatedAt: Date.now() });
await replaceSessionEntry(
{ sessionKey: writerSessionKey, storePath },
{ sessionId: "sqlite-large-archive-writer", updatedAt: Date.now() },
);
await replaceTranscriptEvents({ sessionKey, sessionId, storePath }, events);
await replaceSessionEntry(
{ sessionKey: warmupSessionKey, storePath },
{ sessionId: warmupSessionId, updatedAt: Date.now() },
);
await replaceTranscriptEvents(
{ sessionKey: warmupSessionKey, sessionId: warmupSessionId, storePath },
[
{
type: "session",
id: warmupSessionId,
content: "warm the built archive worker",
} as unknown as TestTranscriptEvent,
],
);
const databasePath = requireSqliteDatabasePath(storePath);
expect(readSessionRowCounts(databasePath, sessionId)).toEqual({
fts: 1,
sessionWindows: 1,
transcriptEvents: events.length,
});
closeOpenClawAgentDatabasesForTest();
await expect(inst.entrypoint()).resolves.toEqual(
expect.arrayContaining([expect.stringMatching(/^dist\/index\.(?:js|mjs)$/u)]),
);
await inst.startGateway();
[deleteClient, probeClient] = await Promise.all([
connectGatewayClient({
url: inst.url,
token: inst.gatewayToken,
clientDisplayName: "sqlite-large-archive-delete",
requestTimeoutMs: 120_000,
timeoutMs: 20_000,
}),
connectGatewayClient({
url: inst.url,
token: inst.gatewayToken,
clientDisplayName: "sqlite-large-archive-presence",
requestTimeoutMs: 2_000,
timeoutMs: 20_000,
}),
]);
// Cold-opening and indexing the pre-seeded 64 MiB database is outside
// the deletion latency measurement below and can exceed the normal RPC
// timeout on Windows CI hosts. Finish that one-time initialization first.
await deleteClient.request("sessions.list", {}, { timeoutMs: 120_000 });
for (let attempt = 0; attempt < 3; attempt += 1) {
await probeClient.request("system-presence", {}, { timeoutMs: 2_000 });
}
// Prime the built sidecar and OS file cache with a tiny transcript so the
// latency assertion below measures data-size-dependent archive work.
await deleteClient.request(
"sessions.delete",
{ key: warmupSessionKey, deleteTranscript: true },
{ timeoutMs: 20_000 },
);
let archivePublishedAt: number | undefined;
let deleteSettled = false;
const publicationPoll = setInterval(() => {
if (findPublishedArchive(archiveDirectory, sessionId)) {
archivePublishedAt ??= performance.now();
}
}, 5);
const deletion = deleteClient
.request<{ archived?: string[]; deleted?: boolean; ok?: boolean }>(
"sessions.delete",
{ key: sessionKey, deleteTranscript: true },
{ timeoutMs: 120_000 },
)
.finally(() => {
deleteSettled = true;
});
void deletion.catch(() => undefined);
const writerStartedAt = performance.now();
const writerResult = await probeClient.request<{ key?: string; ok?: boolean }>(
"sessions.patch",
{ key: writerSessionKey, label: "writer-progressed-during-archive" },
{ timeoutMs: 2_000 },
);
const writerLatencyMs = performance.now() - writerStartedAt;
expect(writerResult).toMatchObject({ ok: true, key: writerSessionKey });
expect(writerLatencyMs).toBeLessThan(500);
expect(deleteSettled).toBe(false);
expect(archivePublishedAt).toBeUndefined();
const prePublicationProbeLatencies: number[] = [];
const shouldProbeBeforePublication = () => !deleteSettled && archivePublishedAt === undefined;
try {
while (shouldProbeBeforePublication()) {
const probeStartedAt = performance.now();
await probeClient.request("system-presence", {}, { timeoutMs: 2_000 });
const probeCompletedAt = performance.now();
// Record every probe that started before publication was observed.
// A synchronous implementation can delay this response until after
// publication; dropping that crossing sample would hide the stall.
prePublicationProbeLatencies.push(probeCompletedAt - probeStartedAt);
await new Promise<void>((resolve) => {
setTimeout(resolve, 5);
});
}
} finally {
clearInterval(publicationPoll);
}
const deleteResult = await deletion;
expect(deleteResult).toMatchObject({ ok: true, deleted: true });
expect(prePublicationProbeLatencies.length).toBeGreaterThan(5);
// Keep enough headroom for Windows scheduling and a probe that crosses
// into the existing synchronous SQLite/FTS deletion tail. The former
// synchronous archive path instead exceeds the probe's 2s RPC timeout.
expect(Math.max(...prePublicationProbeLatencies)).toBeLessThan(500);
const archivedPath = deleteResult.archived?.[0];
expect(archivedPath).toBeTruthy();
await Promise.all([
disconnectGatewayClient(deleteClient),
disconnectGatewayClient(probeClient),
]);
deleteClient = undefined;
probeClient = undefined;
await inst.stopGateway();
const archivedContent = readSessionArchiveContentSync(archivedPath ?? "");
expect(Buffer.byteLength(archivedContent)).toBe(Buffer.byteLength(expectedArchiveContent));
expect(sha256(archivedContent)).toBe(sha256(expectedArchiveContent));
expect(loadSessionEntry({ sessionKey, storePath })).toBeUndefined();
await expect(loadTranscriptEvents({ sessionKey, sessionId, storePath })).resolves.toEqual([]);
expect(readSessionRowCounts(databasePath, sessionId)).toEqual({
fts: 0,
sessionWindows: 0,
transcriptEvents: 0,
});
} finally {
await Promise.allSettled(
[deleteClient, probeClient]
.filter((client): client is NonNullable<typeof client> => client !== undefined)
.map((client) => disconnectGatewayClient(client)),
);
await inst.stopGateway();
closeOpenClawAgentDatabasesForTest();
closeOpenClawStateDatabaseForTest();
await inst.cleanup();
}
}, 180_000);
});
type TestTranscriptEvent = Parameters<typeof replaceTranscriptEvents>[1][number];
function createLargeTranscriptEvents(sessionId: string): TestTranscriptEvent[] {
const indexedMessage = {
type: "message",
id: "sqlite-large-archive-indexed-message",
parentId: null,
message: {
role: "user",
content: [{ type: "text", text: "large archive searchable marker" }],
},
timestamp: Date.now(),
} as unknown as TestTranscriptEvent;
return [
indexedMessage,
...Array.from(
{ length: 63 },
(_, index) =>
({
type: "session",
id: `${sessionId}-${index}`,
content: `${index}:${randomBytes(768 * 1024).toString("base64")}`,
}) as unknown as TestTranscriptEvent,
),
];
}
function findPublishedArchive(archiveDirectory: string, sessionId: string): string | undefined {
const prefix = `${sessionId}.jsonl.deleted.`;
try {
return fs
.readdirSync(archiveDirectory)
.find((entry) => entry.startsWith(prefix) && !entry.endsWith(".tmp"));
} catch {
return undefined;
}
}
function requireSqliteDatabasePath(storePath: string): string {
const target = resolveSqliteTargetFromSessionStorePath(storePath);
if (!target.path) {
throw new Error(`could not resolve SQLite database path for ${storePath}`);
}
return target.path;
}
function readSessionRowCounts(
databasePath: string,
sessionId: string,
): {
fts: number;
sessionWindows: number;
transcriptEvents: number;
} {
const database = new DatabaseSync(databasePath, { readOnly: true });
try {
const count = (table: "session_transcript_fts" | "session_windows" | "transcript_events") => {
const row = database
.prepare(`SELECT COUNT(*) AS count FROM ${table} WHERE session_id = ?`)
.get(sessionId) as { count: number };
return row.count;
};
return {
fts: count("session_transcript_fts"),
sessionWindows: count("session_windows"),
transcriptEvents: count("transcript_events"),
};
} finally {
database.close();
}
}
function sha256(content: string): string {
return createHash("sha256").update(content).digest("hex");
}
@@ -1,5 +1,5 @@
// SQLite sessions/transcripts flip proof test runs the script-style gateway lifecycle probe.
import { describe, expect, it } from "vitest";
import { describe, it } from "vitest";
import { assertSqliteFlipProofCore } from "../helpers/sqlite-sessions-transcripts-flip-proof-assertions.ts";
import { runSqliteSessionsTranscriptsFlipProof } from "../helpers/sqlite-sessions-transcripts-flip-proof.ts";
@@ -8,227 +8,5 @@ describe("SQLite sessions/transcripts flip proof harness", () => {
const report = await runSqliteSessionsTranscriptsFlipProof();
assertSqliteFlipProofCore(report);
expect(report.checkpoints.map((checkpoint) => checkpoint.label)).toEqual([
"seeded-legacy-store",
"after-startup-import",
"after-doctor-inspect",
"after-doctor-validate",
"after-rollback-restore",
"after-gateway-restart",
"after-chat-send",
"after-full-agent-turn",
"after-manual-compaction",
"after-plugin-sdk-consumer",
"after-cleanup-pruning",
"after-doctor-import-idempotence",
"after-downgrade-reupgrade-import",
"after-sqlite-busy-contention",
"after-concurrent-multi-client",
"after-sessions-reset",
"after-second-startup-after-reset",
"after-transcript-append",
"after-sessions-delete",
"after-shared-first-delete",
"after-shared-final-delete",
"after-final-doctor-inspect",
]);
const startupImportCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-startup-import",
);
expect(
startupImportCheckpoint?.archiveArtifacts.some(
(artifact) =>
artifact.path.includes("old-orphan.deleted.jsonl") &&
artifact.textTail?.includes("old-orphan") === true,
),
).toBe(true);
expect(report.rollbackRestore).toMatchObject({
archivedBeforeRestore: true,
failedManifestIssueCode: "e2e_forced_post_archive_failure",
sourceRestored: true,
sqliteStillExists: true,
});
expect(report.rollbackRestore?.manifestPath).toContain("session-sqlite-migration-runs");
expect(
report.rollbackRestore?.restoredFiles.some((filePath) =>
filePath.replaceAll("\\", "/").endsWith("/sqlite-rollback-restore.jsonl"),
),
).toBe(true);
expect(
report.rollbackRestore?.idempotentRestoreSkippedFiles.some((filePath) =>
filePath.replaceAll("\\", "/").endsWith("/sqlite-rollback-restore.jsonl"),
),
).toBe(true);
expect(report.scaleMigration).toMatchObject({
minTranscriptEventsPerSession: 4,
seededEvents: 96,
seededSessions: 24,
});
expect(report.scaleMigration?.importedSessionKeys).toHaveLength(24);
expect(report.scaleMigration?.startupImportElapsedMs).toBeGreaterThanOrEqual(0);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-full-agent-turn" &&
checkpoint.sqlite.trackedEntries.some(
(entry) =>
entry.sessionKey === report.fullTurnSessionKey &&
entry.transcriptEvents >= 2 &&
entry.trajectoryEvents >= 1,
),
),
).toBe(true);
expect(report.manualCompaction).toMatchObject({
checkpointCount: 1,
compacted: true,
sessionKey: report.manualCompactionSessionKey,
});
expect(report.manualCompaction?.transcriptIdentity).toBe(report.manualCompactionSessionKey);
expect(report.manualCompaction?.rowCountBefore).toBeGreaterThanOrEqual(2);
expect(report.manualCompaction?.rowCountAfter).toBeGreaterThanOrEqual(1);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-manual-compaction" &&
checkpoint.activeJsonl.length === 0 &&
checkpoint.sqlite.trackedEntries.some(
(entry) =>
entry.sessionKey === report.manualCompactionSessionKey &&
Array.isArray(entry.entry?.compactionCheckpoints) &&
entry.entry.compactionCheckpoints.length >= 1,
),
),
).toBe(true);
expect(report.pluginSdkConsumer).toMatchObject({
activeTrajectoryPointerForSessionExists: false,
activeTrajectoryRuntimeSidecarForSessionExists: false,
activeTrajectorySessionSidecarForSessionExists: false,
});
expect(report.pluginSdkConsumer?.sessionIdentity).toBe(report.pluginSdkSessionKey);
expect(report.pluginSdkConsumer?.listedSessionKeys).toContain(report.pluginSdkSessionKey);
expect(report.pluginSdkConsumer?.transcriptEventsAfterAppend).toBeGreaterThan(
report.pluginSdkConsumer?.transcriptEventsBeforeAppend ?? 0,
);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-plugin-sdk-consumer" &&
checkpoint.sqlite.trajectoryRuntimeEvents >= 1 &&
checkpoint.sqlite.trackedEntries.some(
(entry) =>
entry.sessionKey === report.pluginSdkSessionKey &&
entry.trajectoryEvents >= 1 &&
entry.transcriptEvents >= 3,
),
),
).toBe(true);
expect(report.downgradeReupgrade).toMatchObject({
activeJsonlArchived: true,
doctorImportedEntries: 1,
doctorImportedTranscriptEvents: 2,
sessionId: "sqlite-downgrade-reupgrade",
sessionKey: "agent:main:dashboard:sqlite-downgrade-reupgrade",
trajectoryPointerArchived: true,
trajectoryPointerSourceRemoved: true,
trajectorySidecarArchived: true,
trajectorySidecarSourceRemoved: true,
transcriptEvents: 2,
});
const downgradeCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-downgrade-reupgrade-import",
);
expect(
downgradeCheckpoint?.archiveArtifacts.some(
(artifact) =>
artifact.path.includes("sqlite-downgrade-reupgrade.trajectory.jsonl") &&
artifact.textTail?.includes("trajectory") === true,
),
).toBe(true);
expect(
downgradeCheckpoint?.archiveArtifacts.some((artifact) =>
artifact.path.includes("sqlite-downgrade-reupgrade.trajectory-path.json"),
),
).toBe(true);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-downgrade-reupgrade-import" &&
checkpoint.sqlite.trackedEntries.some(
(entry) =>
entry.sessionKey === "agent:main:dashboard:sqlite-downgrade-reupgrade" &&
entry.transcriptEvents === 2,
),
),
).toBe(true);
expect(report.busyContention).toMatchObject({
childExitCode: 0,
childSignal: null,
holdMs: 500,
sessionId: "sqlite-busy-contention",
sessionKey: "agent:main:dashboard:sqlite-busy-contention",
transcriptEvents: 2,
});
expect(report.busyContention?.elapsedMs).toBeGreaterThanOrEqual(250);
expect(report.secondStartupAfterReset).toMatchObject({
activeJsonlForSessionExists: false,
historyContainsPostResetAppend: true,
sessionKey: report.resetSessionKey,
});
expect(report.secondStartupAfterReset?.transcriptEvents).toBeGreaterThanOrEqual(1);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-transcript-append" &&
checkpoint.sqlite.trackedEntries.some(
(entry) => entry.sessionKey === report.resetSessionKey && entry.transcriptEvents >= 1,
),
),
).toBe(true);
const deleteCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-sessions-delete",
);
const deleteArchive = deleteCheckpoint?.archiveArtifacts.find(
(artifact) =>
artifact.archiveReason === "deleted" &&
artifact.archiveSessionId === "sqlite-delete-session",
);
expect(deleteArchive?.messageTexts).toContain("delete me");
const sharedFinalCheckpoint = report.checkpoints.find(
(checkpoint) => checkpoint.label === "after-shared-final-delete",
);
const sharedFinalArchive = sharedFinalCheckpoint?.archiveArtifacts.find(
(artifact) =>
artifact.archiveReason === "deleted" &&
artifact.archiveSessionId === "sqlite-shared-session",
);
const retainedSharedImportSources = sharedFinalCheckpoint?.archiveArtifacts.filter(
(artifact) =>
artifact.path.includes("session-sqlite-import-archive") &&
(artifact.path.includes("sqlite-shared-a.jsonl") ||
artifact.path.includes("sqlite-shared-b.jsonl")),
);
expect(
sharedFinalArchive?.messageTexts?.includes("shared") ||
(retainedSharedImportSources?.length === 2 &&
retainedSharedImportSources.every((artifact) =>
artifact.messageTexts?.some((text) => text.includes("shared")),
)),
).toBe(true);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-shared-first-delete" &&
checkpoint.sqlite.trackedEntries.some(
(entry) => entry.sessionKey === report.sharedSessionKeys[1],
),
),
).toBe(true);
expect(
report.checkpoints.some(
(checkpoint) =>
checkpoint.label === "after-shared-final-delete" &&
checkpoint.archiveArtifacts.length > 0,
),
).toBe(true);
}, 420_000);
});