From e5cee36b46cce2e1cbaa857eaa9963d141dee8d3 Mon Sep 17 00:00:00 2001 From: WhatsSkiLL Date: Thu, 30 Jul 2026 11:47:59 +0200 Subject: [PATCH] fix(memory): retry failed queued session targets (#115923) * fix(memory): retain failed queued sync targets * fix(memory): drain retained targets on idle sync * test(memory): prove idle queued sync recovery * fix(memory): preserve queued sync ownership * fix(memory): avoid queued sync self-deadlock * fix(memory): stop queued recovery during close * fix(memory): clear retained sync state on close * test(memory): prove live queued rejection transition * fix(memory): enforce sync repro invariants * test(memory): bound archive recovery proof * test(memory): seed archive proof transcript * test(memory): exercise archived transcript recovery * chore(knip): register memory sync repro * fix(memory): reject blank queries before settings * style(memory): satisfy queue recovery lint * test(memory): align doctor migration expectations * test(memory): insert explicit provenance fixture --------- Co-authored-by: IWhatsskill <284122573+IWhatsskill@users.noreply.github.com> Co-authored-by: Peter Steinberger --- config/knip.config.ts | 1 + .../memory-core/src/memory/index.test.ts | 504 +++++++++++++++++- .../src/memory/manager-sync-control.ts | 62 ++- .../src/memory/manager-targeted-sync.test.ts | 157 +++++- .../memory/manager.readonly-recovery.test.ts | 8 + extensions/memory-core/src/memory/manager.ts | 42 +- scripts/memory-index-manager.sync-repro.ts | 347 ++++++++++++ 7 files changed, 1108 insertions(+), 13 deletions(-) create mode 100644 scripts/memory-index-manager.sync-repro.ts diff --git a/config/knip.config.ts b/config/knip.config.ts index fea9325546da..ce6ba73167f5 100644 --- a/config/knip.config.ts +++ b/config/knip.config.ts @@ -65,6 +65,7 @@ const repositoryScriptEntries = [ // Invoked by scripts/lib/live-docker-stage.sh during container validation. "scripts/live-docker-normalize-config.ts!", "scripts/mcp-code-mode-gateway-e2e.ts!", + "scripts/memory-index-manager.sync-repro.ts!", "scripts/openclaw-release-clawhub-plan.ts!", "scripts/openclaw-release-clawhub-runtime-state.ts!", // Oxlint loads this JS plugin by path from config/oxlint/boundary-guards.json. diff --git a/extensions/memory-core/src/memory/index.test.ts b/extensions/memory-core/src/memory/index.test.ts index b30c5b74969e..170cd50e95c2 100644 --- a/extensions/memory-core/src/memory/index.test.ts +++ b/extensions/memory-core/src/memory/index.test.ts @@ -3,12 +3,14 @@ import { mkdirSync, rmSync } from "node:fs"; 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 { clearMemoryEmbeddingProviders as clearRegistry } from "openclaw/plugin-sdk/memory-core-host-engine-embeddings"; import { hashText, INVALID_PROJECT_ANNOTATION_KEY, MEMORY_CHUNKING_VERSION, + type MemorySessionSyncTarget, + type MemorySyncParams, } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import { resolveSessionTranscriptsDirForAgent } from "openclaw/plugin-sdk/memory-core-host-runtime-core"; import { upsertSessionEntry } from "openclaw/plugin-sdk/session-store-runtime"; @@ -1733,6 +1735,506 @@ describe("memory index", () => { } }); + it("drains retained queued targets through the next idle sync call", async () => { + const markers = { + blocker: "BLOCKER LOCKED SYNC 729", + retained: "RETAINED RETRY TARGET 729", + trigger: "IDLE TRIGGER TARGET 729", + }; + const sessionKey = (sessionId: string) => `agent:main:proof:${sessionId}`; + const manager = await getFreshManager( + createCfg({ + provider: "none", + sources: ["sessions"], + sessionMemory: true, + }), + ); + let lock: DatabaseSync | null = null; + try { + await manager.sync({ reason: "test-baseline", force: true }); + for (const [sessionId, marker] of Object.entries(markers)) { + await seedMemoryIndexSessionTranscript({ + sessionId, + sessionKey: sessionKey(sessionId), + messages: [ + { + role: "user", + timestamp: Date.now(), + content: marker, + }, + ], + }); + } + + const dbPath = resolveOpenClawAgentSqlitePath({ agentId: "main" }); + lock = new DatabaseSync(dbPath); + lock.exec("PRAGMA busy_timeout = 0"); + lock.exec("BEGIN EXCLUSIVE"); + + const active = manager.sync({ + reason: "test-locked-owner", + sessions: [ + { + agentId: "main", + sessionId: "blocker", + sessionKey: sessionKey("blocker"), + }, + ], + }); + const failedQueued = manager.sync({ + reason: "test-queued-retained", + sessions: [ + { + agentId: "main", + sessionId: "retained", + sessionKey: sessionKey("retained"), + }, + ], + }); + const failures = await Promise.allSettled([active, failedQueued]); + lock.exec("ROLLBACK"); + lock.close(); + lock = null; + const describeSqliteFailure = (failure: unknown): string => { + const details = [String(failure)]; + if (failure && typeof failure === "object") { + const record = failure as Record; + for (const key of ["message", "code"] as const) { + if (typeof record[key] === "string") { + details.push(record[key]); + } + } + if (record.cause && typeof record.cause === "object") { + const cause = record.cause as Record; + for (const key of ["message", "code"] as const) { + if (typeof cause[key] === "string") { + details.push(cause[key]); + } + } + } + } + return details.join(" "); + }; + for (const result of failures) { + expect(result.status).toBe("rejected"); + if (result.status !== "rejected") { + throw new Error("expected SQLite-locked sync to reject"); + } + expect(describeSqliteFailure(result.reason)).toMatch( + /SQLITE_(?:BUSY|LOCKED)|database is (?:busy|locked)/i, + ); + } + + const ftsMatchCount = (marker: string): number => { + const observer = new DatabaseSync(dbPath, { readOnly: true }); + try { + return ( + observer + .prepare( + "SELECT COUNT(*) AS count FROM memory_index_chunks_fts WHERE memory_index_chunks_fts MATCH ?", + ) + .get(`"${marker}"`) as { count: number } + ).count; + } finally { + observer.close(); + } + }; + + expect(ftsMatchCount(markers.retained)).toBe(0); + expect(ftsMatchCount(markers.trigger)).toBe(0); + const recoveryState = manager as unknown as { + syncing: Promise | null; + queuedSessions: Map; + sessionsDirtyFiles: Set; + sessionsFullRetryDirty: boolean; + }; + expect(recoveryState.syncing).toBeNull(); + expect(recoveryState.queuedSessions.size).toBe(1); + expect(recoveryState.sessionsDirtyFiles.size).toBe(0); + expect(recoveryState.sessionsFullRetryDirty).toBe(false); + + const recoveryProgress = vi.fn(); + const recovery = manager.sync({ + reason: "test-recovery-trigger", + sessions: [ + { + agentId: "main", + sessionId: "trigger", + sessionKey: sessionKey("trigger"), + }, + ], + progress: recoveryProgress, + }); + // A full sync can claim `syncing` before the retained queue owner resumes. + // Both owners must settle without the queue awaiting its own promise. + const competingFullSync = manager.sync({ reason: "test-competing-full-sync" }); + const recoveryResults = await Promise.allSettled([recovery, competingFullSync]); + expect(recoveryResults.map((result) => result.status)).toEqual(["fulfilled", "fulfilled"]); + + expect(ftsMatchCount(markers.retained)).toBeGreaterThan(0); + expect(ftsMatchCount(markers.trigger)).toBeGreaterThan(0); + expect(recoveryState.queuedSessions.size).toBe(0); + expect(recoveryProgress).toHaveBeenCalled(); + } finally { + if (lock) { + try { + lock.exec("ROLLBACK"); + } finally { + lock.close(); + } + } + await manager.close?.(); + } + }); + + it("drains retained queued targets from a live rejection transition", async () => { + const markers = { + retained: "LIVE REJECTION RETAINED TARGET 729", + transition: "LIVE REJECTION TRANSITION TARGET 729", + trigger: "LIVE REJECTION RECOVERY TARGET 729", + }; + const sessionKey = (sessionId: string) => `agent:main:live-rejection:${sessionId}`; + const manager = await getFreshManager( + createCfg({ + provider: "none", + sources: ["sessions"], + sessionMemory: true, + }), + ); + let resolveActiveSync: (() => void) | undefined; + const activeSyncGate = new Promise((resolve) => { + resolveActiveSync = resolve; + }); + let rejectQueuedSync: ((error: Error) => void) | undefined; + const queuedSyncGate = new Promise((_resolve, reject) => { + rejectQueuedSync = reject; + }); + const owner = manager as unknown as { + syncing: Promise | null; + queuedSessions: Map; + queuedSessionSync: Promise | null; + runSyncWithReadonlyRecovery: (params?: MemorySyncParams) => Promise; + }; + const runSyncWithReadonlyRecovery = owner.runSyncWithReadonlyRecovery.bind(owner); + const runSync = vi + .spyOn(owner, "runSyncWithReadonlyRecovery") + .mockImplementationOnce(async (params) => await runSyncWithReadonlyRecovery(params)) + .mockImplementationOnce(async () => await activeSyncGate) + .mockImplementationOnce(async () => await queuedSyncGate) + .mockImplementation(async (params) => await runSyncWithReadonlyRecovery(params)); + const queuedError = new Error("controlled queued rejection"); + try { + await manager.sync({ reason: "test-live-rejection-baseline", force: true }); + for (const [sessionId, marker] of Object.entries(markers)) { + await seedMemoryIndexSessionTranscript({ + sessionId, + sessionKey: sessionKey(sessionId), + messages: [ + { + role: "user", + timestamp: Date.now(), + content: marker, + }, + ], + }); + } + + const active = manager.sync({ + reason: "test-live-rejection-owner", + sessions: [ + { + agentId: "main", + sessionId: "active", + sessionKey: sessionKey("active"), + }, + ], + }); + const queuedProgress = vi.fn(); + const failedQueued = manager.sync({ + reason: "test-live-rejection-queued", + sessions: [ + { + agentId: "main", + sessionId: "retained", + sessionKey: sessionKey("retained"), + }, + ], + force: true, + progress: queuedProgress, + }); + const failuresPromise = Promise.allSettled([active, failedQueued]); + resolveActiveSync?.(); + await vi.waitFor(() => { + expect(runSync).toHaveBeenCalledTimes(3); + expect(owner.syncing).not.toBeNull(); + expect(owner.queuedSessionSync).not.toBeNull(); + }); + const rejectingQueuedSync = owner.syncing; + if (!rejectingQueuedSync) { + throw new Error("expected a live queued sync"); + } + + let resolveTransitionResult!: (result: PromiseSettledResult) => void; + const transitionResult = new Promise>((resolve) => { + resolveTransitionResult = resolve; + }); + let transitionState: + | { syncingNull: boolean; queueOwnerLive: boolean; queuedTargets: number } + | undefined; + const transitionProgress = vi.fn(); + void rejectingQueuedSync.catch(() => { + transitionState = { + syncingNull: owner.syncing === null, + queueOwnerLive: owner.queuedSessionSync !== null, + queuedTargets: owner.queuedSessions.size, + }; + const transitionCall = manager.sync({ + reason: "test-live-rejection-transition", + sessions: [ + { + agentId: "main", + sessionId: "transition", + sessionKey: sessionKey("transition"), + }, + ], + progress: transitionProgress, + }); + void transitionCall.then( + (value) => resolveTransitionResult({ status: "fulfilled", value }), + (reason: unknown) => resolveTransitionResult({ status: "rejected", reason }), + ); + }); + + rejectQueuedSync?.(queuedError); + const failures = await failuresPromise; + const transitionFailure = await transitionResult; + expect(failures[0]?.status).toBe("fulfilled"); + expect(failures[1]?.status).toBe("rejected"); + expect(transitionFailure.status).toBe("rejected"); + if (failures[1]?.status !== "rejected" || transitionFailure.status !== "rejected") { + throw new Error("expected shared queued rejection"); + } + expect(failures[1].reason).toBe(queuedError); + expect(transitionFailure.reason).toBe(queuedError); + expect(transitionState).toEqual({ + syncingNull: true, + queueOwnerLive: true, + queuedTargets: 0, + }); + expect(Array.from(owner.queuedSessions.values())).toEqual([ + { + agentId: "main", + sessionId: "transition", + sessionKey: sessionKey("transition"), + }, + { + agentId: "main", + sessionId: "retained", + sessionKey: sessionKey("retained"), + }, + ]); + expect(queuedProgress).not.toHaveBeenCalled(); + expect(transitionProgress).not.toHaveBeenCalled(); + + const recoveryProgress = vi.fn(); + await manager.sync({ + reason: "test-live-rejection-recovery", + sessions: [ + { + agentId: "main", + sessionId: "trigger", + sessionKey: sessionKey("trigger"), + }, + ], + progress: recoveryProgress, + }); + + const dbPath = resolveOpenClawAgentSqlitePath({ agentId: "main" }); + const observer = new DatabaseSync(dbPath, { readOnly: true }); + try { + const indexedCount = (marker: string) => + ( + observer + .prepare("SELECT COUNT(*) AS count FROM memory_index_chunks WHERE text LIKE ?") + .get(`%${marker}%`) as { count: number } + ).count; + expect(indexedCount(markers.retained)).toBeGreaterThan(0); + expect(indexedCount(markers.transition)).toBeGreaterThan(0); + expect(indexedCount(markers.trigger)).toBeGreaterThan(0); + } finally { + observer.close(); + } + expect(owner.queuedSessions.size).toBe(0); + expect(recoveryProgress).toHaveBeenCalled(); + expect(transitionProgress).not.toHaveBeenCalled(); + } finally { + resolveActiveSync?.(); + rejectQueuedSync?.(queuedError); + await manager.close?.(); + runSync.mockRestore(); + } + }); + + it("clears retained queued targets when close interrupts a competing sync", async () => { + const manager = await getFreshManager( + createCfg({ + provider: "none", + sources: ["sessions"], + sessionMemory: true, + }), + ); + let resolveFullSync: (() => void) | undefined; + const fullSyncGate = new Promise((resolve) => { + resolveFullSync = resolve; + }); + const owner = manager as unknown as { + closing: boolean; + closed: boolean; + queuedSessions: Map; + queuedProgressCallbacks: Set>; + queuedForce: boolean; + syncAdmitted: (params?: MemorySyncParams) => Promise; + runSyncWithReadonlyRecovery: (params?: MemorySyncParams) => Promise; + }; + const syncAdmitted = vi.spyOn(owner, "syncAdmitted"); + const runSyncWithReadonlyRecovery = vi + .spyOn(owner, "runSyncWithReadonlyRecovery") + .mockReturnValueOnce(fullSyncGate); + const progress = vi.fn(); + owner.queuedSessions.set("retained", { + agentId: "main", + sessionId: "retained-close", + sessionKey: "agent:main:retained-close", + }); + + try { + const recovery = manager.sync({ + reason: "test-close-recovery", + sessions: [ + { + agentId: "main", + sessionId: "trigger-close", + sessionKey: "agent:main:trigger-close", + }, + ], + force: true, + progress, + }); + const competingFullSync = manager.sync({ reason: "test-close-competing-full-sync" }); + + await vi.waitFor(() => { + expect(syncAdmitted).toHaveBeenCalledTimes(2); + }); + const closing = manager.close?.() ?? Promise.resolve(); + expect(owner.closing).toBe(true); + resolveFullSync?.(); + + await expect(Promise.all([recovery, competingFullSync, closing])).resolves.toEqual([ + undefined, + undefined, + undefined, + ]); + expect(runSyncWithReadonlyRecovery).toHaveBeenCalledTimes(1); + expect(syncAdmitted).toHaveBeenCalledTimes(2); + expect(owner.closed).toBe(true); + expect(owner.queuedSessions.size).toBe(0); + expect(owner.queuedProgressCallbacks.size).toBe(0); + expect(owner.queuedForce).toBe(false); + expect(progress).not.toHaveBeenCalled(); + } finally { + resolveFullSync?.(); + await manager.close?.(); + runSyncWithReadonlyRecovery.mockRestore(); + syncAdmitted.mockRestore(); + } + }); + + it("clears retained queued targets after failure when the manager closes", async () => { + const manager = await getFreshManager( + createCfg({ + provider: "none", + sources: ["sessions"], + sessionMemory: true, + }), + ); + let resolveActiveSync: (() => void) | undefined; + const activeSyncGate = new Promise((resolve) => { + resolveActiveSync = resolve; + }); + const owner = manager as unknown as { + closed: boolean; + queuedArchiveFiles: Set; + queuedSessions: Map; + queuedProgressCallbacks: Set>; + queuedForce: boolean; + queuedSessionSync: Promise | null; + runSyncWithReadonlyRecovery: (params?: MemorySyncParams) => Promise; + }; + const runSyncWithReadonlyRecovery = vi + .spyOn(owner, "runSyncWithReadonlyRecovery") + .mockReturnValueOnce(activeSyncGate) + .mockRejectedValueOnce(new Error("test queued failure")); + const progress = vi.fn(); + + try { + const active = manager.sync({ + reason: "test-close-after-failure-owner", + sessions: [ + { + agentId: "main", + sessionId: "active-close-after-failure", + sessionKey: "agent:main:active-close-after-failure", + }, + ], + }); + const failedQueued = manager.sync({ + reason: "test-close-after-failure-queued", + sessions: [ + { + agentId: "main", + sessionId: "retained-close-after-failure", + sessionKey: "agent:main:retained-close-after-failure", + }, + ], + archiveFiles: ["/tmp/retained-close-after-failure.jsonl"], + force: true, + progress, + }); + const queuedRejection = expect(failedQueued).rejects.toThrow("test queued failure"); + + resolveActiveSync?.(); + await active; + await queuedRejection; + + expect(runSyncWithReadonlyRecovery).toHaveBeenCalledTimes(2); + expect(owner.queuedArchiveFiles).toEqual( + new Set(["/tmp/retained-close-after-failure.jsonl"]), + ); + expect(Array.from(owner.queuedSessions.values())).toEqual([ + { + agentId: "main", + sessionId: "retained-close-after-failure", + sessionKey: "agent:main:retained-close-after-failure", + }, + ]); + expect(owner.queuedForce).toBe(true); + expect(owner.queuedProgressCallbacks.size).toBe(0); + expect(owner.queuedSessionSync).toBeNull(); + + await manager.close?.(); + + expect(owner.closed).toBe(true); + expect(owner.queuedArchiveFiles.size).toBe(0); + expect(owner.queuedSessions.size).toBe(0); + expect(owner.queuedProgressCallbacks.size).toBe(0); + expect(owner.queuedForce).toBe(false); + } finally { + resolveActiveSync?.(); + await manager.close?.(); + runSyncWithReadonlyRecovery.mockRestore(); + } + }); + it("keeps provider cutover vector search paused during targeted session sync", async () => { try { setMemoryIndexStateDir(path.join(workspaceDir, ".state-targeted-cutover")); diff --git a/extensions/memory-core/src/memory/manager-sync-control.ts b/extensions/memory-core/src/memory/manager-sync-control.ts index db40e90f2cc8..56a9eb4b83cd 100644 --- a/extensions/memory-core/src/memory/manager-sync-control.ts +++ b/extensions/memory-core/src/memory/manager-sync-control.ts @@ -122,11 +122,14 @@ export function enqueueMemoryTargetedSessionSync( getSyncing: () => Promise | null; getQueuedArchiveFiles: () => Set; getQueuedSessions: () => Map; + getQueuedForce: () => boolean; + setQueuedForce: (value: boolean) => void; + getQueuedProgressCallbacks: () => Set>; getQueuedSessionSync: () => Promise | null; setQueuedSessionSync: (value: Promise | null) => void; sync: (params?: MemorySyncParams) => Promise; }, - targets?: Pick, + targets?: Pick, ): Promise { const queuedArchiveFiles = state.getQueuedArchiveFiles(); for (const sessionFile of targets?.archiveFiles ?? []) { @@ -145,6 +148,12 @@ export function enqueueMemoryTargetedSessionSync( if (queuedArchiveFiles.size === 0 && queuedSessions.size === 0) { return state.getSyncing() ?? Promise.resolve(); } + if (targets?.force) { + state.setQueuedForce(true); + } + if (targets?.progress) { + state.getQueuedProgressCallbacks().add(targets.progress); + } if (!state.getQueuedSessionSync()) { state.setQueuedSessionSync( (async () => { @@ -156,15 +165,56 @@ export function enqueueMemoryTargetedSessionSync( ) { const pendingArchiveFiles = Array.from(state.getQueuedArchiveFiles()); const pendingSessions = Array.from(state.getQueuedSessions().values()); + const pendingForce = state.getQueuedForce(); + const pendingProgressCallbacks = Array.from(state.getQueuedProgressCallbacks()); state.getQueuedArchiveFiles().clear(); state.getQueuedSessions().clear(); - await state.sync({ - reason: "queued-sessions", - sessions: pendingSessions, - archiveFiles: pendingArchiveFiles, - }); + state.setQueuedForce(false); + state.getQueuedProgressCallbacks().clear(); + const progress = + pendingProgressCallbacks.length > 0 + ? (update: MemorySyncProgressUpdate) => { + for (const callback of pendingProgressCallbacks) { + callback(update); + } + } + : undefined; + try { + await state.sync({ + reason: "queued-sessions", + ...(pendingForce ? { force: true } : {}), + sessions: pendingSessions, + archiveFiles: pendingArchiveFiles, + ...(progress ? { progress } : {}), + }); + } catch (err) { + // Merge the failed batch with arrivals queued during sync so the + // next trigger can retry every target instead of dropping work. + for (const archiveFile of pendingArchiveFiles) { + state.getQueuedArchiveFiles().add(archiveFile); + } + for (const session of pendingSessions) { + state.getQueuedSessions().set(memorySessionSyncTargetKey(session), session); + } + if (pendingForce) { + state.setQueuedForce(true); + } + // Every caller awaiting this queue owner receives the rejection. + // Do not retain callbacks that could otherwise fire after their + // originating promise has already failed. + state.getQueuedProgressCallbacks().clear(); + throw err; + } } } finally { + if (state.isClosed()) { + // A closed manager cannot drain retained work. Release every + // manager-owned target and caller closure with the queue owner. + state.getQueuedArchiveFiles().clear(); + state.getQueuedSessions().clear(); + state.setQueuedForce(false); + state.getQueuedProgressCallbacks().clear(); + } state.setQueuedSessionSync(null); } })(), diff --git a/extensions/memory-core/src/memory/manager-targeted-sync.test.ts b/extensions/memory-core/src/memory/manager-targeted-sync.test.ts index 9e3dc295377b..917323ea86b1 100644 --- a/extensions/memory-core/src/memory/manager-targeted-sync.test.ts +++ b/extensions/memory-core/src/memory/manager-targeted-sync.test.ts @@ -1,5 +1,8 @@ // Memory Core tests cover manager targeted sync plugin behavior. -import type { MemorySessionSyncTarget } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; +import type { + MemorySessionSyncTarget, + MemorySyncParams, +} from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import { describe, expect, it, vi } from "vitest"; import { enqueueMemoryTargetedSessionSync } from "./manager-sync-control.js"; import { @@ -80,8 +83,14 @@ describe("memory targeted session sync", () => { }); const queuedArchiveFiles = new Set(); const queuedSessions = new Map(); + let queuedForce = false; + const queuedProgressCallbacks = new Set>(); let queuedSessionSync: Promise | null = null; - const sync = vi.fn(async () => {}); + const progressUpdate = { completed: 1, total: 2, label: "queued" }; + const progress = vi.fn(); + const sync = vi.fn(async (params?: MemorySyncParams) => { + params?.progress?.(progressUpdate); + }); const queued = enqueueMemoryTargetedSessionSync( { @@ -89,6 +98,11 @@ describe("memory targeted session sync", () => { getSyncing: () => syncing, getQueuedArchiveFiles: () => queuedArchiveFiles, getQueuedSessions: () => queuedSessions, + getQueuedForce: () => queuedForce, + setQueuedForce: (value) => { + queuedForce = value; + }, + getQueuedProgressCallbacks: () => queuedProgressCallbacks, getQueuedSessionSync: () => queuedSessionSync, setQueuedSessionSync: (value) => { queuedSessionSync = value; @@ -97,6 +111,8 @@ describe("memory targeted session sync", () => { }, { sessions: [{ agentId: "main", sessionId: "targeted", sessionKey: "agent:main:targeted" }], + force: true, + progress, }, ); @@ -105,8 +121,145 @@ describe("memory targeted session sync", () => { expect(sync).toHaveBeenCalledWith({ reason: "queued-sessions", + force: true, sessions: [{ agentId: "main", sessionId: "targeted", sessionKey: "agent:main:targeted" }], archiveFiles: [], + progress: expect.any(Function), }); + expect(progress).toHaveBeenCalledWith(progressUpdate); + }); + + it("keeps failed queued targets for a later retry", async () => { + let resolveSyncing: (() => void) | undefined; + const syncing = new Promise((resolve) => { + resolveSyncing = resolve; + }); + let rejectQueuedSync: ((error: Error) => void) | undefined; + const queuedSync = new Promise((_resolve, reject) => { + rejectQueuedSync = reject; + }); + const queuedArchiveFiles = new Set(); + const queuedSessions = new Map(); + let queuedForce = false; + const queuedProgressCallbacks = new Set>(); + let queuedSessionSync: Promise | null = null; + const sync = vi.fn().mockReturnValueOnce(queuedSync).mockResolvedValueOnce(undefined); + const state = { + isClosed: () => false, + getSyncing: () => syncing, + getQueuedArchiveFiles: () => queuedArchiveFiles, + getQueuedSessions: () => queuedSessions, + getQueuedForce: () => queuedForce, + setQueuedForce: (value: boolean) => { + queuedForce = value; + }, + getQueuedProgressCallbacks: () => queuedProgressCallbacks, + getQueuedSessionSync: () => queuedSessionSync, + setQueuedSessionSync: (value: Promise | null) => { + queuedSessionSync = value; + }, + sync, + }; + + const firstProgress = vi.fn(); + const first = enqueueMemoryTargetedSessionSync(state, { + sessions: [{ agentId: "main", sessionId: "first", sessionKey: "agent:main:first" }], + archiveFiles: ["/tmp/first.jsonl"], + force: true, + progress: firstProgress, + }); + const firstRejection = expect(first).rejects.toThrow("transient sqlite failure"); + + resolveSyncing?.(); + await vi.waitFor(() => { + expect(sync).toHaveBeenCalledTimes(1); + }); + + const concurrentProgress = vi.fn(); + const concurrent = enqueueMemoryTargetedSessionSync(state, { + sessions: [{ agentId: "main", sessionId: "second", sessionKey: "agent:main:second" }], + archiveFiles: ["/tmp/second.jsonl"], + progress: concurrentProgress, + }); + expect(concurrent).toBe(first); + + rejectQueuedSync?.(new Error("transient sqlite failure")); + await firstRejection; + + expect(queuedArchiveFiles).toEqual(new Set(["/tmp/second.jsonl", "/tmp/first.jsonl"])); + expect(Array.from(queuedSessions.values())).toEqual([ + { agentId: "main", sessionId: "second", sessionKey: "agent:main:second" }, + { agentId: "main", sessionId: "first", sessionKey: "agent:main:first" }, + ]); + expect(queuedSessionSync).toBeNull(); + expect(queuedProgressCallbacks.size).toBe(0); + expect(queuedForce).toBe(true); + + await enqueueMemoryTargetedSessionSync(state); + + expect(sync).toHaveBeenCalledTimes(2); + expect(sync).toHaveBeenLastCalledWith({ + reason: "queued-sessions", + force: true, + sessions: [ + { agentId: "main", sessionId: "second", sessionKey: "agent:main:second" }, + { agentId: "main", sessionId: "first", sessionKey: "agent:main:first" }, + ], + archiveFiles: ["/tmp/second.jsonl", "/tmp/first.jsonl"], + }); + expect(queuedArchiveFiles.size).toBe(0); + expect(queuedSessions.size).toBe(0); + expect(queuedSessionSync).toBeNull(); + expect(queuedForce).toBe(false); + }); + + it("clears queued state when the manager closes while the queue waits", async () => { + let resolveSyncing: (() => void) | undefined; + const syncing = new Promise((resolve) => { + resolveSyncing = resolve; + }); + let closed = false; + const queuedArchiveFiles = new Set(["/tmp/close-retained.jsonl"]); + const queuedSessions = new Map(); + let queuedForce = false; + const queuedProgressCallbacks = new Set>(); + let queuedSessionSync: Promise | null = null; + const sync = vi.fn(async () => undefined); + const progress = vi.fn(); + + const queued = enqueueMemoryTargetedSessionSync( + { + isClosed: () => closed, + getSyncing: () => syncing, + getQueuedArchiveFiles: () => queuedArchiveFiles, + getQueuedSessions: () => queuedSessions, + getQueuedForce: () => queuedForce, + setQueuedForce: (value) => { + queuedForce = value; + }, + getQueuedProgressCallbacks: () => queuedProgressCallbacks, + getQueuedSessionSync: () => queuedSessionSync, + setQueuedSessionSync: (value) => { + queuedSessionSync = value; + }, + sync, + }, + { + sessions: [{ agentId: "main", sessionId: "close", sessionKey: "agent:main:close" }], + force: true, + progress, + }, + ); + + closed = true; + resolveSyncing?.(); + await queued; + + expect(sync).not.toHaveBeenCalled(); + expect(queuedArchiveFiles.size).toBe(0); + expect(queuedSessions.size).toBe(0); + expect(queuedProgressCallbacks.size).toBe(0); + expect(queuedForce).toBe(false); + expect(queuedSessionSync).toBeNull(); }); }); diff --git a/extensions/memory-core/src/memory/manager.readonly-recovery.test.ts b/extensions/memory-core/src/memory/manager.readonly-recovery.test.ts index a54e1e5b87fb..7cc12f0aeb21 100644 --- a/extensions/memory-core/src/memory/manager.readonly-recovery.test.ts +++ b/extensions/memory-core/src/memory/manager.readonly-recovery.test.ts @@ -3,6 +3,7 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import type { DatabaseSync } from "node:sqlite"; +import type { MemorySyncParams } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { closeMemoryDatabase, openMemoryDatabaseAtPath } from "./manager-db.js"; import { @@ -34,6 +35,8 @@ describe("memory manager readonly recovery", () => { function createQueuedSyncHarness(syncing: Promise) { const queuedArchiveFiles = new Set(); const queuedSessions = new Map(); + let queuedForce = false; + const queuedProgressCallbacks = new Set>(); let queuedSessionSync: Promise | null = null; const sync = vi.fn(async () => {}); return { @@ -48,6 +51,11 @@ describe("memory manager readonly recovery", () => { getSyncing: () => syncing, getQueuedArchiveFiles: () => queuedArchiveFiles, getQueuedSessions: () => queuedSessions, + getQueuedForce: () => queuedForce, + setQueuedForce: (value: boolean) => { + queuedForce = value; + }, + getQueuedProgressCallbacks: () => queuedProgressCallbacks, getQueuedSessionSync: () => queuedSessionSync, setQueuedSessionSync: (value: Promise | null) => { queuedSessionSync = value; diff --git a/extensions/memory-core/src/memory/manager.ts b/extensions/memory-core/src/memory/manager.ts index 3865f6885e7d..62ecb65a7a0e 100644 --- a/extensions/memory-core/src/memory/manager.ts +++ b/extensions/memory-core/src/memory/manager.ts @@ -491,6 +491,8 @@ export class MemoryIndexManager extends MemoryManagerEmbeddingOps implements Mem private syncing: Promise | null = null; private queuedArchiveFiles = new Set(); private queuedSessions = new Map(); + private queuedForce = false; + private queuedProgressCallbacks = new Set>(); private queuedSessionSync: Promise | null = null; private readonlyRecoveryAttempts = 0; private readonlyRecoverySuccesses = 0; @@ -1895,15 +1897,38 @@ export class MemoryIndexManager extends MemoryManagerEmbeddingOps implements Mem if (this.closing || this.closed) { return; } + if ( + hasTargetedSessionSyncParams(params) && + (this.queuedSessionSync !== null || + this.queuedArchiveFiles.size > 0 || + this.queuedSessions.size > 0) + ) { + // A failed queued batch stays manager-owned. Route the next targeted + // call through the queue even while idle so it adopts that retained work. + return await this.enqueueTargetedSessionSync(params); + } return await this.syncAdmitted(params); } private async syncAdmitted( params?: MemorySyncParams, - options?: { allowEmbeddingBootstrapFallback?: boolean }, + options?: { + allowEmbeddingBootstrapFallback?: boolean; + queuedSessionOwner?: boolean; + }, ): Promise { if (this.syncing) { if (hasTargetedSessionSyncParams(params)) { + if (options?.queuedSessionOwner) { + // Another caller claimed the sync slot after this queue owner was + // created. Wait for it, then retry admission instead of enqueueing + // into the promise that is already awaiting this call. + await this.syncing.catch(() => undefined); + if (this.closing || this.closed) { + return; + } + return await this.syncAdmitted(params, options); + } return this.enqueueTargetedSessionSync(params); } try { @@ -1995,19 +2020,24 @@ export class MemoryIndexManager extends MemoryManagerEmbeddingOps implements Mem } private enqueueTargetedSessionSync( - targets?: Pick, + targets?: Pick, ): Promise { return enqueueMemoryTargetedSessionSync( { - isClosed: () => this.closed, + isClosed: () => this.closing || this.closed, getSyncing: () => this.syncing, getQueuedArchiveFiles: () => this.queuedArchiveFiles, getQueuedSessions: () => this.queuedSessions, + getQueuedForce: () => this.queuedForce, + setQueuedForce: (value) => { + this.queuedForce = value; + }, + getQueuedProgressCallbacks: () => this.queuedProgressCallbacks, getQueuedSessionSync: () => this.queuedSessionSync, setQueuedSessionSync: (value) => { this.queuedSessionSync = value; }, - sync: async (params) => await this.syncAdmitted(params), + sync: async (params) => await this.syncAdmitted(params, { queuedSessionOwner: true }), }, targets, ); @@ -2314,6 +2344,10 @@ export class MemoryIndexManager extends MemoryManagerEmbeddingOps implements Mem private async closeOnce(): Promise { this.closing = true; + this.queuedArchiveFiles.clear(); + this.queuedSessions.clear(); + this.queuedForce = false; + this.queuedProgressCallbacks.clear(); await this.awaitManagerIdle(); this.closed = true; const pendingProviderInit = this.providerInitPromise; diff --git a/scripts/memory-index-manager.sync-repro.ts b/scripts/memory-index-manager.sync-repro.ts new file mode 100644 index 000000000000..861ecdb0cb48 --- /dev/null +++ b/scripts/memory-index-manager.sync-repro.ts @@ -0,0 +1,347 @@ +import { execFileSync } from "node:child_process"; +import fs from "node:fs/promises"; +import path from "node:path"; +import { DatabaseSync } from "node:sqlite"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/memory-core-host-engine-foundation"; +import { resolveSessionTranscriptsDirForAgent } from "openclaw/plugin-sdk/memory-core-host-runtime-core"; +import { upsertSessionEntry } from "openclaw/plugin-sdk/session-store-runtime"; +import { appendSessionTranscriptMessageByIdentity } from "openclaw/plugin-sdk/session-transcript-runtime"; +import { resolveOpenClawAgentSqlitePath } from "openclaw/plugin-sdk/sqlite-runtime"; +import { + closeAllMemorySearchManagers, + getMemorySearchManager, +} from "../extensions/memory-core/src/memory/index.ts"; + +const proofRoot = process.argv[2]; +const exactHead = process.argv[3]; +if (!proofRoot || !exactHead?.match(/^[0-9a-f]{40}$/)) { + throw new Error("proof root and exact 40-character head are required"); +} +const observedHead = execFileSync("git", ["rev-parse", "HEAD"], { + encoding: "utf8", +}).trim(); +if (observedHead !== exactHead) { + throw new Error(`exact head mismatch: expected ${exactHead}, observed ${observedHead}`); +} +const dirtyState = execFileSync("git", ["status", "--porcelain"], { + encoding: "utf8", +}).trim(); +if (dirtyState) { + throw new Error("the production repro requires a clean exact-head worktree"); +} + +const agentId = "main"; +const stateDir = path.join(proofRoot, "state"); +const workspaceDir = path.join(proofRoot, "workspace"); +const markers = { + blocker: "BLOCKER_LOCKED_SYNC_729", + retained: "RETAINED_RETRY_TARGET_729", + trigger: "CONCURRENT_TRIGGER_TARGET_729", + archive: "RETAINED_ARCHIVE_TARGET_729", +}; + +Reflect.set(process.env, "OPENCLAW_STATE_DIR", stateDir); +await fs.mkdir(path.join(workspaceDir, "memory"), { recursive: true }); +await fs.writeFile(path.join(workspaceDir, "MEMORY.md"), "# Proof workspace\n"); + +const cfg: OpenClawConfig = { + memory: { + search: { + provider: "none", + sources: ["sessions"], + rememberAcrossConversations: true, + store: { vector: { enabled: false } }, + query: { minScore: 0 }, + }, + }, + agents: { + defaults: { workspace: workspaceDir }, + list: [{ id: agentId, default: true }], + }, +}; + +async function seedSession(sessionId: string, marker: string): Promise { + const sessionsDir = resolveSessionTranscriptsDirForAgent(agentId); + const storePath = path.join(sessionsDir, "sessions.json"); + const sessionKey = `agent:${agentId}:proof:${sessionId}`; + await fs.mkdir(sessionsDir, { recursive: true }); + await upsertSessionEntry({ + agentId, + sessionKey, + storePath, + entry: { sessionId, updatedAt: Date.now() }, + }); + await appendSessionTranscriptMessageByIdentity({ + agentId, + sessionId, + sessionKey, + storePath, + message: { + role: "user", + timestamp: Date.now(), + content: [{ type: "text", text: marker }], + }, + }); + return sessionKey; +} + +function openExclusiveLock(dbPath: string): DatabaseSync { + const db = new DatabaseSync(dbPath); + db.exec("PRAGMA busy_timeout = 0"); + db.exec("BEGIN EXCLUSIVE"); + return db; +} + +function releaseExclusiveLock(db: DatabaseSync | null): void { + if (!db) { + return; + } + try { + db.exec("ROLLBACK"); + } finally { + db.close(); + } +} + +function describeSqliteFailure(failure: unknown): string { + const details = [String(failure)]; + if (failure && typeof failure === "object") { + const record = failure as Record; + for (const key of ["message", "code"] as const) { + if (typeof record[key] === "string") { + details.push(record[key]); + } + } + if (record.cause && typeof record.cause === "object") { + const cause = record.cause as Record; + for (const key of ["message", "code"] as const) { + if (typeof cause[key] === "string") { + details.push(cause[key]); + } + } + } + } + return details.join(" "); +} + +function isSqliteLockFailure(failure: unknown): boolean { + return /SQLITE_(?:BUSY|LOCKED)|database is (?:busy|locked)/i.test(describeSqliteFailure(failure)); +} + +async function withTimeout(promise: Promise, timeoutMs: number, label: string): Promise { + let timer: ReturnType | undefined; + const timeout = new Promise((_resolve, reject) => { + timer = setTimeout(() => { + reject(new Error(`${label} timed out after ${timeoutMs}ms`)); + }, timeoutMs); + }); + try { + return await Promise.race([promise, timeout]); + } finally { + if (timer) { + clearTimeout(timer); + } + } +} + +let lock: DatabaseSync | null = null; +try { + const result = await getMemorySearchManager({ cfg, agentId }); + if (!result.manager) { + throw new Error(`memory manager unavailable: ${result.error ?? "unknown"}`); + } + const manager = result.manager; + const sync = manager.sync?.bind(manager); + if (!sync) { + throw new Error("memory manager sync is unavailable"); + } + await sync({ reason: "proof-baseline", force: true }); + + const blockerKey = await seedSession("proof-blocker", markers.blocker); + const retainedKey = await seedSession("proof-retained", markers.retained); + const triggerKey = await seedSession("proof-trigger", markers.trigger); + const archiveFile = path.join( + resolveSessionTranscriptsDirForAgent(agentId), + "proof-archive.jsonl.deleted.2026-07-29T00-00-00.000Z", + ); + await fs.writeFile( + archiveFile, + [ + JSON.stringify({ + type: "session", + id: "proof-archive", + timestamp: new Date().toISOString(), + }), + JSON.stringify({ + type: "message", + message: { + role: "user", + timestamp: Date.now(), + content: [{ type: "text", text: markers.archive }], + }, + }), + ].join("\n") + "\n", + "utf8", + ); + const dbPath = resolveOpenClawAgentSqlitePath({ agentId }); + + lock = openExclusiveLock(dbPath); + const blockedOwner = sync({ + reason: "proof-locked-owner", + sessions: [{ agentId, sessionId: "proof-blocker", sessionKey: blockerKey }], + }); + const failedQueued = sync({ + reason: "proof-queued-retained", + sessions: [{ agentId, sessionId: "proof-retained", sessionKey: retainedKey }], + archiveFiles: [archiveFile], + }); + const failures = await Promise.allSettled([blockedOwner, failedQueued]); + const lockedSyncFailures = failures.filter((entry) => entry.status === "rejected").length; + const sqliteLockFailures = failures.filter( + (entry) => entry.status === "rejected" && isSqliteLockFailure(entry.reason), + ).length; + releaseExclusiveLock(lock); + lock = null; + if (lockedSyncFailures !== 2) { + throw new Error(`expected two locked sync failures, received ${lockedSyncFailures}`); + } + if (sqliteLockFailures !== 2) { + throw new Error(`expected two SQLite lock failures, received ${sqliteLockFailures}`); + } + + const observer = new DatabaseSync(dbPath, { readOnly: true }); + const retainedBefore = + ( + observer + .prepare("SELECT COUNT(*) AS count FROM memory_index_chunks WHERE text LIKE ?") + .get(`%${markers.retained}%`) as { count: number } + ).count > 0; + const triggerBefore = + ( + observer + .prepare("SELECT COUNT(*) AS count FROM memory_index_chunks WHERE text LIKE ?") + .get(`%${markers.trigger}%`) as { count: number } + ).count > 0; + const archiveBefore = + ( + observer + .prepare("SELECT COUNT(*) AS count FROM memory_index_chunks WHERE text LIKE ?") + .get(`%${markers.archive}%`) as { count: number } + ).count > 0; + observer.close(); + if (retainedBefore || triggerBefore || archiveBefore) { + throw new Error("a recovery target was unexpectedly indexed before the idle trigger"); + } + + const recoveryState = manager as unknown as { + syncing: Promise | null; + queuedArchiveFiles: Set; + queuedSessions: Map; + sessionsDirtyFiles: Set; + sessionsFullRetryDirty: boolean; + }; + const retainedQueueBeforeRecovery = recoveryState.queuedSessions.size; + const retainedArchiveQueueBeforeRecovery = recoveryState.queuedArchiveFiles.size; + const dirtySessionFilesBeforeRecovery = recoveryState.sessionsDirtyFiles.size; + const fullRetryBeforeRecovery = recoveryState.sessionsFullRetryDirty; + if ( + recoveryState.syncing !== null || + retainedQueueBeforeRecovery !== 1 || + retainedArchiveQueueBeforeRecovery !== 1 + ) { + throw new Error("manager was not idle with exactly one retained session and archive target"); + } + + const recoveryProgress: Array<{ completed: number; total: number; label?: string }> = []; + const recovery = sync({ + reason: "proof-idle-recovery-trigger", + sessions: [{ agentId, sessionId: "proof-trigger", sessionKey: triggerKey }], + progress: (update) => recoveryProgress.push(update), + }); + // Start an untargeted sync before the retained queue owner resumes. + // Its distinct public progress callback proves that this call, rather than + // an implementation token, reached competing production admission. + const competingUntargetedProgress: Array<{ completed: number; total: number; label?: string }> = + []; + const competingUntargetedSync = sync({ + reason: "proof-competing-untargeted-sync", + progress: (update) => competingUntargetedProgress.push(update), + }); + const queueSettlementTimeoutMs = 15_000; + const recoveryResults = await withTimeout( + Promise.allSettled([recovery, competingUntargetedSync]), + queueSettlementTimeoutMs, + "queue-owner self-deadlock check", + ); + const recoveryStatus = recoveryResults[0]?.status; + const competingUntargetedStatus = recoveryResults[1]?.status; + if (recoveryStatus !== "fulfilled" || competingUntargetedStatus !== "fulfilled") { + throw new Error( + `concurrent recovery did not settle: recovery=${recoveryStatus ?? "missing"} untargeted=${competingUntargetedStatus ?? "missing"}`, + ); + } + if (competingUntargetedProgress.length === 0) { + throw new Error("competing untargeted sync emitted no public progress updates"); + } + + const recoveryObserver = new DatabaseSync(dbPath, { readOnly: true }); + const indexedCount = (marker: string) => + ( + recoveryObserver + .prepare("SELECT COUNT(*) AS count FROM memory_index_chunks WHERE text LIKE ?") + .get(`%${marker}%`) as { count: number } + ).count; + const retainedAfter = indexedCount(markers.retained) > 0; + const triggerAfter = indexedCount(markers.trigger) > 0; + const archiveAfter = indexedCount(markers.archive) > 0; + recoveryObserver.close(); + if (!retainedAfter || !triggerAfter || !archiveAfter) { + throw new Error( + `recovery result mismatch: retained=${String(retainedAfter)} trigger=${String(triggerAfter)} archive=${String(archiveAfter)}`, + ); + } + if (recoveryProgress.length === 0) { + throw new Error("idle recovery trigger did not receive progress"); + } + const retainedQueueAfterRecovery = recoveryState.queuedSessions.size; + const retainedArchiveQueueAfterRecovery = recoveryState.queuedArchiveFiles.size; + if ( + dirtySessionFilesBeforeRecovery !== 0 || + fullRetryBeforeRecovery || + retainedQueueAfterRecovery !== 0 || + retainedArchiveQueueAfterRecovery !== 0 + ) { + throw new Error( + `unexpected recovery ownership state: dirty=${dirtySessionFilesBeforeRecovery} fullRetry=${String(fullRetryBeforeRecovery)} retainedAfter=${retainedQueueAfterRecovery} archiveAfter=${retainedArchiveQueueAfterRecovery}`, + ); + } + + console.log(`exact_head=${exactHead}`); + console.log("test_runner=none"); + console.log("entrypoint=MemoryIndexManager.sync"); + console.log("owners=memory-manager,session-store,sqlite"); + console.log(`locked_sync_failures=${lockedSyncFailures}`); + console.log("locked_sync_failure_kind=sqlite-busy"); + console.log("recovery_manager_state=idle"); + console.log("recovery_input_sessions=proof-trigger"); + console.log(`recovery_progress_updates=${recoveryProgress.length}`); + console.log(`competing_untargeted_sync_progress_updates=${competingUntargetedProgress.length}`); + console.log(`recovery_sync_status=${recoveryStatus}`); + console.log(`competing_untargeted_sync_status=${competingUntargetedStatus}`); + console.log(`queue_settlement_timeout_ms=${queueSettlementTimeoutMs}`); + console.log(`retained_queue_before_recovery=${retainedQueueBeforeRecovery}`); + console.log(`retained_archive_queue_before_recovery=${retainedArchiveQueueBeforeRecovery}`); + console.log(`sessions_dirty_files_before_recovery=${dirtySessionFilesBeforeRecovery}`); + console.log(`sessions_full_retry_dirty_before_recovery=${String(fullRetryBeforeRecovery)}`); + console.log("retained_target_before_recovery=absent"); + console.log("retained_archive_target_before_recovery=absent"); + console.log("retained_target_after_recovery=indexed"); + console.log("retained_archive_target_after_recovery=indexed"); + console.log("idle_trigger_after_recovery=indexed"); + console.log(`retained_queue_after_recovery=${retainedQueueAfterRecovery}`); + console.log(`retained_archive_queue_after_recovery=${retainedArchiveQueueAfterRecovery}`); + console.log("verdict=pass"); +} finally { + releaseExclusiveLock(lock); + await closeAllMemorySearchManagers(); +}