From 8616acbfd407a5746c0288f0e8a56b37eebb314a Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 1 Aug 2026 11:02:42 -0700 Subject: [PATCH] fix(memory): recover restored session freshness --- .../memory/manager-session-sync-state.test.ts | 8 + .../src/memory/manager-session-sync-state.ts | 4 +- .../src/memory/manager-source-sync-ops.ts | 27 +- .../manager-sync-ops.startup-catchup.test.ts | 265 +++++++++++++++++- 4 files changed, 289 insertions(+), 15 deletions(-) diff --git a/extensions/memory-core/src/memory/manager-session-sync-state.test.ts b/extensions/memory-core/src/memory/manager-session-sync-state.test.ts index 80d674014a42..fd0846619efa 100644 --- a/extensions/memory-core/src/memory/manager-session-sync-state.test.ts +++ b/extensions/memory-core/src/memory/manager-session-sync-state.test.ts @@ -106,6 +106,12 @@ describe("memory session sync state", () => { mtimeMs: 250, size: 20, }, + { + absPath: "/tmp/sessions/rolled-back.jsonl", + path: "sessions/rolled-back.jsonl", + mtimeMs: 150, + size: 20, + }, { absPath: "/tmp/sessions/resized.jsonl", path: "sessions/resized.jsonl", @@ -124,6 +130,7 @@ describe("memory session sync state", () => { { path: "sessions/sub-ms-newer.jsonl", hash: "hash-sub-ms", mtime: 100.25, size: 10 }, { path: "sessions/invalidated.jsonl", hash: "", mtime: 200, size: 20 }, { path: "sessions/newer.jsonl", hash: "hash-newer", mtime: 200, size: 20 }, + { path: "sessions/rolled-back.jsonl", hash: "hash-rolled-back", mtime: 200, size: 20 }, { path: "sessions/resized.jsonl", hash: "hash-resized", mtime: 300, size: 30 }, ], }); @@ -132,6 +139,7 @@ describe("memory session sync state", () => { "/tmp/sessions/sub-ms-newer.jsonl", "/tmp/sessions/invalidated.jsonl", "/tmp/sessions/newer.jsonl", + "/tmp/sessions/rolled-back.jsonl", "/tmp/sessions/resized.jsonl", "/tmp/sessions/missing.jsonl", ]); diff --git a/extensions/memory-core/src/memory/manager-session-sync-state.ts b/extensions/memory-core/src/memory/manager-session-sync-state.ts index 7857762185c6..7be696a96892 100644 --- a/extensions/memory-core/src/memory/manager-session-sync-state.ts +++ b/extensions/memory-core/src/memory/manager-session-sync-state.ts @@ -26,7 +26,9 @@ export function resolveMemorySessionStartupDirtyFiles(params: { dirtyFiles.push(file.absPath); continue; } - if (file.size !== indexedSize || file.mtimeMs > indexedMtimeMs) { + // File mtimes and SQLite session updatedAt values can move backward after + // restore/reset. The downstream content-hash gate suppresses unchanged rewrites. + if (file.size !== indexedSize || file.mtimeMs !== indexedMtimeMs) { dirtyFiles.push(file.absPath); } } diff --git a/extensions/memory-core/src/memory/manager-source-sync-ops.ts b/extensions/memory-core/src/memory/manager-source-sync-ops.ts index 2d9bcad63a27..0f3b992a97fc 100644 --- a/extensions/memory-core/src/memory/manager-source-sync-ops.ts +++ b/extensions/memory-core/src/memory/manager-source-sync-ops.ts @@ -174,6 +174,29 @@ export abstract class MemoryManagerSourceSyncOps extends MemoryManagerSessionSyn const deleteChunksByPathAndSource = this.db.prepare( `DELETE FROM memory_index_chunks WHERE path = ? AND source = ?`, ); + const updateUnchangedSessionSourceMetadata = this.db.prepare( + `UPDATE memory_index_sources + SET mtime = ?, size = ? + WHERE path = ? AND source = 'sessions' AND hash = ?`, + ); + const refreshUnchangedSessionSourceMetadata = (entry: MemoryIndexEntry): boolean => { + // Hash equality preserves chunks and embeddings; only converge the source + // fingerprint so restored sessions do not repeat catch-up on every startup. + return ( + updateUnchangedSessionSourceMetadata.run(entry.mtimeMs, entry.size, entry.path, entry.hash) + .changes === 1 + ); + }; + const canSkipUnchangedSessionEntry = ( + entry: MemoryIndexEntry, + absPath: string, + existingHash: string | undefined, + ): boolean => { + if (params.needsFullReindex || existingHash !== entry.hash) { + return false; + } + return !this.sessionsDirtyFiles.has(absPath) || refreshUnchangedSessionSourceMetadata(entry); + }; const deleteFtsRowsByPathAndSource = this.fts.enabled && this.fts.available ? this.db.prepare(`DELETE FROM ${FTS_TABLE} WHERE path = ? AND source = ?`) @@ -340,7 +363,7 @@ export abstract class MemoryManagerSourceSyncOps extends MemoryManagerSessionSyn path: entry.path, existingHashes, }); - if (!params.needsFullReindex && existingHash === entry.hash) { + if (canSkipUnchangedSessionEntry(entry, absPath, existingHash)) { if (params.progress) { params.progress.completed += 1; params.progress.report({ @@ -412,7 +435,7 @@ export abstract class MemoryManagerSourceSyncOps extends MemoryManagerSessionSyn path: entry.path, existingHashes, }); - if (!params.needsFullReindex && existingHash === entry.hash) { + if (canSkipUnchangedSessionEntry(entry, absPath, existingHash)) { if (params.progress) { params.progress.completed += 1; params.progress.report({ diff --git a/extensions/memory-core/src/memory/manager-sync-ops.startup-catchup.test.ts b/extensions/memory-core/src/memory/manager-sync-ops.startup-catchup.test.ts index a1324deae121..968c2d8891d5 100644 --- a/extensions/memory-core/src/memory/manager-sync-ops.startup-catchup.test.ts +++ b/extensions/memory-core/src/memory/manager-sync-ops.startup-catchup.test.ts @@ -2,13 +2,16 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; -import type { DatabaseSync } from "node:sqlite"; +import { DatabaseSync } from "node:sqlite"; import { resolveSessionTranscriptsDirForAgent, type OpenClawConfig, type ResolvedMemorySearchConfig, } from "openclaw/plugin-sdk/memory-core-host-engine-foundation"; -import { statSessionEntrySync } from "openclaw/plugin-sdk/memory-core-host-engine-qmd"; +import { + buildSessionEntry, + statSessionEntrySync, +} from "openclaw/plugin-sdk/memory-core-host-engine-qmd"; import { MEMORY_CHUNKING_VERSION, type MemorySource, @@ -64,8 +67,43 @@ type MemorySessionTranscriptUpdate = { const originalStartupStateDir = process.env.OPENCLAW_STATE_DIR; const originalStartupConfigPath = process.env.OPENCLAW_CONFIG_PATH; let transcriptUpdateListener: ((update: MemorySessionTranscriptUpdate) => void) | undefined; +const startupHarnessDatabases = new Set(); type SourceStateRow = { path: string; hash: string; mtime: number; size: number }; + +function createStartupHarnessDatabase(sourceRows: SourceStateRow[]): DatabaseSync { + const db = new DatabaseSync(":memory:"); + db.exec(` + CREATE TABLE memory_index_sources ( + path TEXT NOT NULL, + source TEXT NOT NULL, + hash TEXT NOT NULL, + mtime REAL NOT NULL, + size INTEGER NOT NULL, + UNIQUE(path, source) + ); + CREATE TABLE memory_index_chunks ( + id TEXT PRIMARY KEY, + path TEXT NOT NULL, + source TEXT NOT NULL, + model TEXT NOT NULL + ); + CREATE TABLE memory_index_source_update_audit (path TEXT NOT NULL); + CREATE TRIGGER memory_index_source_update_audit_trigger + AFTER UPDATE ON memory_index_sources + BEGIN + INSERT INTO memory_index_source_update_audit (path) VALUES (NEW.path); + END; + `); + const insert = db.prepare( + `INSERT INTO memory_index_sources (path, source, hash, mtime, size) VALUES (?, 'sessions', ?, ?, ?)`, + ); + for (const row of sourceRows) { + insert.run(row.path, row.hash, row.mtime, row.size); + } + startupHarnessDatabases.add(db); + return db; +} function setStartupStateDir(stateDir: string): void { Reflect.set(process.env, "OPENCLAW_STATE_DIR", stateDir); } @@ -148,16 +186,37 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps { sourceRows: SourceStateRow[], private readonly indexSessionUpdates = false, private readonly subscribeToRealEvents = false, + private readonly deferSessionIndex = false, + database?: DatabaseSync, ) { super(); this.sources.add("sessions"); - this.db = { - prepare: () => ({ - all: () => sourceRows, - get: () => undefined, - run: () => undefined, - }), - } as unknown as DatabaseSync; + this.db = database ?? createStartupHarnessDatabase(sourceRows); + } + + restartForStartup(): SessionStartupCatchupHarness { + return new SessionStartupCatchupHarness( + [], + this.indexSessionUpdates, + false, + this.deferSessionIndex, + this.db, + ); + } + + getIndexedSourceState(pathname: string): SourceStateRow | undefined { + return this.db + .prepare( + `SELECT path, hash, mtime, size FROM memory_index_sources WHERE path = ? AND source = 'sessions'`, + ) + .get(pathname) as SourceStateRow | undefined; + } + + getSourceMetadataUpdateCount(): number { + const row = this.db + .prepare(`SELECT COUNT(*) AS count FROM memory_index_source_update_audit`) + .get() as { count: number }; + return row.count; } async catchUp(): Promise { @@ -172,6 +231,13 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps { await this.runSync(params); } + async runArchiveSyncForTest(): Promise { + await this.syncArchiveFiles({ + needsFullReindex: false, + deferIndex: this.deferSessionIndex, + }); + } + getDirtyArchiveFiles(): string[] { return Array.from(this.sessionsDirtyFiles); } @@ -273,7 +339,10 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps { protected async sync(params?: MemorySyncParams): Promise { this.syncCalls.push(params ?? {}); this.pendingSyncWork = this.indexSessionUpdates - ? this.syncArchiveFiles({ needsFullReindex: false }).then(() => undefined) + ? this.syncArchiveFiles({ + needsFullReindex: false, + deferIndex: this.deferSessionIndex, + }).then(() => undefined) : Promise.resolve(); await this.pendingSyncWork; } @@ -333,20 +402,29 @@ describe("session startup catch-up", () => { restoreStartupEnv(); clearRuntimeConfigSnapshot(); clearConfigCache(); + for (const database of startupHarnessDatabases) { + database.close(); + } + startupHarnessDatabases.clear(); closeOpenClawAgentDatabasesForTest(); await fs.rm(stateDir, { recursive: true, force: true }); }); async function writeSessionFile( name: string, + content = "startup catchup", + timestamp?: string, ): Promise<{ filePath: string; size: number; mtimeMs: number }> { const sessionsDir = resolveSessionTranscriptsDirForAgent("main"); await fs.mkdir(sessionsDir, { recursive: true }); const filePath = path.join(sessionsDir, name); await fs.writeFile( filePath, - JSON.stringify({ type: "message", message: { role: "user", content: "startup catchup" } }) + - "\n", + JSON.stringify({ + type: "message", + ...(timestamp ? { timestamp } : {}), + message: { role: "user", content }, + }) + "\n", "utf-8", ); const stat = await fs.stat(filePath); @@ -533,6 +611,169 @@ describe("session startup catch-up", () => { expect(harness.syncCalls).toEqual([]); }); + it("indexes a same-size file transcript whose mtime rolled back", async () => { + const archiveName = "thread.jsonl.deleted.2026-08-01T10-00-00.000Z"; + const original = await writeSessionFile(archiveName, "version before"); + const originalEntry = await buildSessionEntry(original.filePath); + if (!originalEntry) { + throw new Error("expected original file transcript entry"); + } + const replacement = await writeSessionFile(archiveName, "version after!"); + expect(replacement.size).toBe(original.size); + const rolledBackMtime = new Date(Math.max(1, original.mtimeMs - 60_000)); + await fs.utimes(replacement.filePath, rolledBackMtime, rolledBackMtime); + const rolledBack = await fs.stat(replacement.filePath); + expect(rolledBack.mtimeMs).toBeLessThan(original.mtimeMs); + + const harness = new SessionStartupCatchupHarness( + [ + { + path: originalEntry.path, + hash: originalEntry.hash, + mtime: original.mtimeMs, + size: original.size, + }, + ], + true, + ); + + await expect(harness.catchUp()).resolves.toEqual([replacement.filePath]); + await harness.waitForSessionSync(); + + expect(harness.syncCalls).toEqual([{ reason: "session-startup-catchup" }]); + expect(harness.indexedPaths).toEqual([`sessions/main/${archiveName}`]); + expect(harness.indexedContents).toEqual(["User: version after!"]); + }); + + it("converges an unchanged file mtime rollback after direct session sync", async () => { + const archiveName = "thread.jsonl.deleted.2026-08-01T11-00-00.000Z"; + const messageTimestamp = "2026-08-01T10:30:00.000Z"; + const original = await writeSessionFile(archiveName, "unchanged content", messageTimestamp); + const originalEntry = await buildSessionEntry(original.filePath); + if (!originalEntry) { + throw new Error("expected original file transcript entry"); + } + const rolledBackMtime = new Date(Math.max(1, original.mtimeMs - 60_000)); + await fs.utimes(original.filePath, rolledBackMtime, rolledBackMtime); + const restoredEntry = await buildSessionEntry(original.filePath); + if (!restoredEntry) { + throw new Error("expected restored file transcript entry"); + } + expect(restoredEntry.hash).toBe(originalEntry.hash); + + const harness = new SessionStartupCatchupHarness( + [ + { + path: originalEntry.path, + hash: originalEntry.hash, + mtime: original.mtimeMs, + size: original.size, + }, + ], + true, + ); + + await expect(harness.catchUp()).resolves.toEqual([original.filePath]); + await harness.waitForSessionSync(); + expect(harness.indexedPaths).toEqual([]); + expect(harness.getIndexedSourceState(originalEntry.path)).toEqual({ + path: originalEntry.path, + hash: originalEntry.hash, + mtime: restoredEntry.mtimeMs, + size: restoredEntry.size, + }); + expect(harness.getSourceMetadataUpdateCount()).toBe(1); + + const restarted = harness.restartForStartup(); + await expect(restarted.catchUp()).resolves.toEqual([]); + expect(restarted.syncCalls).toEqual([]); + expect(restarted.indexedPaths).toEqual([]); + }); + + it("indexes a SQLite transcript whose updatedAt rolled back", async () => { + const session = await writeSqliteSession({ + content: "SQLite rollback", + updatedAt: 10, + }); + const state = statSessionEntrySync(session.sessionKey, { + agentId: "main", + sessionId: session.sessionId, + storePath: session.storePath, + sessionKey: session.sessionKey, + updatedAtMs: 10, + }); + if (!state) { + throw new Error("expected SQLite transcript state"); + } + const harness = new SessionStartupCatchupHarness( + [ + { + path: state.path, + hash: "previous-hash", + mtime: 20, + size: state.size, + }, + ], + true, + ); + + await expect(harness.catchUp()).resolves.toEqual([session.sessionKey]); + await harness.waitForSessionSync(); + + expect(harness.syncCalls).toEqual([{ reason: "session-startup-catchup" }]); + expect(harness.indexedPaths).toEqual([session.corpusPath]); + expect(harness.indexedContents).toEqual(["User: SQLite rollback"]); + }); + + it("converges an unchanged SQLite updatedAt rollback after deferred session sync", async () => { + const session = await writeSqliteSession({ updatedAt: 10 }); + const entry = await buildSessionEntry(session.sessionKey, { + agentId: "main", + sessionId: session.sessionId, + storePath: session.storePath, + sessionKey: session.sessionKey, + updatedAtMs: 10, + sessionKind: "interactive", + }); + if (!entry) { + throw new Error("expected SQLite transcript entry"); + } + const harness = new SessionStartupCatchupHarness( + [ + { + path: entry.path, + hash: entry.hash, + mtime: 20, + size: entry.size, + }, + ], + true, + false, + true, + ); + + await expect(harness.catchUp()).resolves.toEqual([session.sessionKey]); + await harness.waitForSessionSync(); + + expect(harness.syncCalls).toEqual([{ reason: "session-startup-catchup" }]); + expect(harness.indexedPaths).toEqual([]); + expect(harness.indexedContents).toEqual([]); + expect(harness.getIndexedSourceState(entry.path)).toEqual({ + path: entry.path, + hash: entry.hash, + mtime: entry.mtimeMs, + size: entry.size, + }); + expect(harness.getSourceMetadataUpdateCount()).toBe(1); + + const restarted = harness.restartForStartup(); + await expect(restarted.catchUp()).resolves.toEqual([]); + expect(restarted.syncCalls).toEqual([]); + expect(restarted.indexedPaths).toEqual([]); + await restarted.runArchiveSyncForTest(); + expect(restarted.getSourceMetadataUpdateCount()).toBe(1); + }); + it("does not fall back to full session sync when identity targets normalize away", async () => { await writeSessionFile("thread.jsonl"); const harness = new SessionStartupCatchupHarness([]);