Skip to content

Commit 239c193

Browse files
committed
fix(memory-core): guard searches during safe reindex
1 parent 1a3ce7c commit 239c193

4 files changed

Lines changed: 380 additions & 25 deletions

File tree

extensions/memory-core/src/memory/manager-sync-ops.ts

Lines changed: 101 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,16 @@ export abstract class MemoryManagerSyncOps {
260260
>();
261261
protected vectorDegradedWriteWarningShown = false;
262262
private lastMetaSerialized: string | null = null;
263+
// A full reindex briefly repoints this.db at a half-built temp DB (see
264+
// performSafeReindex). Concurrent readers must never observe it: searches wait
265+
// on this latch and status() serves a pre-swap snapshot of the durable index.
266+
// Refcounted because a search-triggered fallback reindex can overlap a
267+
// sync-driven one; the snapshot/latch clear only when the last one settles.
268+
protected reindexing: Promise<void> | null = null;
269+
private reindexDepth = 0;
270+
private resolveReindexing: (() => void) | null = null;
271+
private activeIndexReaders = 0;
272+
private activeIndexReaderResolvers: Array<() => void> = [];
263273

264274
protected abstract readonly cache: { enabled: boolean; maxEntries?: number };
265275
protected abstract db: DatabaseSync;
@@ -284,6 +294,60 @@ export abstract class MemoryManagerSyncOps {
284294
options: { source: MemorySource; content?: string },
285295
): Promise<void>;
286296

297+
// Wait for any in-flight full reindex so reads come from the stable swapped-in
298+
// index, never the temporary build DB that performSafeReindex assigns to this.db.
299+
protected async drainReindex(): Promise<void> {
300+
while (this.reindexing) {
301+
// Only wait for the swap to settle; a failed reindex restores the durable
302+
// DB, so reads proceed against it rather than inheriting the reindex error.
303+
try {
304+
await this.reindexing;
305+
} catch {}
306+
}
307+
}
308+
309+
protected async withIndexRead<T>(read: () => Promise<T>): Promise<T> {
310+
while (true) {
311+
await this.drainReindex();
312+
this.activeIndexReaders++;
313+
// drainReindex can yield before the reader is counted; re-check the latch
314+
// so a newly published reindex cannot race past an unregistered reader.
315+
if (!this.reindexing) {
316+
break;
317+
}
318+
this.releaseIndexRead();
319+
}
320+
321+
try {
322+
return await read();
323+
} finally {
324+
this.releaseIndexRead();
325+
}
326+
}
327+
328+
private async drainActiveIndexReaders(): Promise<void> {
329+
while (this.activeIndexReaders > 0) {
330+
await new Promise<void>((resolve) => {
331+
this.activeIndexReaderResolvers.push(resolve);
332+
});
333+
}
334+
}
335+
336+
private releaseIndexRead(): void {
337+
this.activeIndexReaders--;
338+
if (this.activeIndexReaders === 0) {
339+
const resolvers = this.activeIndexReaderResolvers.splice(0);
340+
for (const resolve of resolvers) {
341+
resolve();
342+
}
343+
}
344+
}
345+
346+
// Hooks for the subclass to snapshot/restore reader-facing state (status()) so
347+
// it keeps reflecting the durable index while this.db points at the temp build.
348+
protected captureReindexReadSnapshot(): void {}
349+
protected clearReindexReadSnapshot(): void {}
350+
287351
protected hasIndexedChunks(): boolean {
288352
const row = this.db.prepare(`SELECT 1 as found FROM chunks LIMIT 1`).get() as
289353
| { found?: number }
@@ -2012,7 +2076,35 @@ export abstract class MemoryManagerSyncOps {
20122076
return true;
20132077
}
20142078

2015-
protected async runSafeReindex(params: {
2079+
protected runSafeReindex(params: {
2080+
reason?: string;
2081+
force?: boolean;
2082+
progress?: MemorySyncProgressState;
2083+
}): Promise<void> {
2084+
// Publish the latch and snapshot the durable index for readers before the
2085+
// swap. Overlapping reindexes share the outermost snapshot/latch, which clear
2086+
// only once the last one settles (depth returns to 0).
2087+
if (this.reindexDepth === 0) {
2088+
this.captureReindexReadSnapshot();
2089+
this.reindexing = new Promise<void>((resolve) => {
2090+
this.resolveReindexing = resolve;
2091+
});
2092+
}
2093+
this.reindexDepth++;
2094+
// Starting the body runs synchronously through `this.db = tempDb` before its
2095+
// first await, so the latch is already published when readers can observe temp.
2096+
return this.performSafeReindex(params).finally(() => {
2097+
this.reindexDepth--;
2098+
if (this.reindexDepth === 0) {
2099+
this.clearReindexReadSnapshot();
2100+
this.reindexing = null;
2101+
this.resolveReindexing?.();
2102+
this.resolveReindexing = null;
2103+
}
2104+
});
2105+
}
2106+
2107+
private async performSafeReindex(params: {
20162108
reason?: string;
20172109
force?: boolean;
20182110
progress?: MemorySyncProgressState;
@@ -2051,6 +2143,7 @@ export abstract class MemoryManagerSyncOps {
20512143
this.vectorReady = originalDbClosed ? null : originalState.vectorReady;
20522144
};
20532145

2146+
await this.drainActiveIndexReaders();
20542147
this.db = tempDb;
20552148
this.resetVectorState();
20562149
this.fts.available = false;
@@ -2132,6 +2225,13 @@ export abstract class MemoryManagerSyncOps {
21322225
this.db = openMemoryDatabaseAtPath(dbPath, this.settings.store.vector.enabled);
21332226
this.resetVectorState();
21342227
this.ensureSchema();
2228+
// The temp build wrote this metadata, but the post-swap handle is the
2229+
// canonical reader DB. Re-write idempotently so immediate readers never see
2230+
// chunks from the new index without the matching identity row.
2231+
if (nextMeta) {
2232+
this.lastMetaSerialized = null;
2233+
this.writeMeta(nextMeta);
2234+
}
21352235
this.vector.dims = nextMeta?.vectorDims;
21362236
} catch (err) {
21372237
try {
Lines changed: 219 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,219 @@
1+
import fs from "node:fs/promises";
2+
import os from "node:os";
3+
import path from "node:path";
4+
import type { OpenClawConfig } from "openclaw/plugin-sdk/memory-core-host-engine-foundation";
5+
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
6+
import { closeAllMemorySearchManagers, getMemorySearchManager } from "./index.js";
7+
import type { MemoryIndexManager } from "./manager.js";
8+
import "./test-runtime-mocks.js";
9+
10+
// FTS-only mode: no embedding provider, so reindex needs no network and stays deterministic.
11+
vi.mock("./embeddings.js", () => ({
12+
createEmbeddingProvider: async () => ({
13+
requestedProvider: "auto",
14+
provider: null,
15+
providerUnavailableReason: "No embeddings provider available.",
16+
}),
17+
resolveEmbeddingProviderAdapterId: () => null,
18+
resolveEmbeddingProviderFallbackModel: () => "fts-only",
19+
}));
20+
21+
type Deferred = { promise: Promise<void>; resolve: () => void };
22+
function deferred(): Deferred {
23+
let resolve = () => {};
24+
const promise = new Promise<void>((res) => {
25+
resolve = res;
26+
});
27+
return { promise, resolve };
28+
}
29+
30+
describe("memory manager reindex read race", () => {
31+
let fixtureRoot = "";
32+
let caseId = 0;
33+
let workspaceDir = "";
34+
let indexPath = "";
35+
let manager: MemoryIndexManager | null = null;
36+
37+
beforeAll(async () => {
38+
fixtureRoot = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-mem-reindex-race-"));
39+
});
40+
41+
beforeEach(async () => {
42+
workspaceDir = path.join(fixtureRoot, `case-${caseId++}`);
43+
await fs.mkdir(path.join(workspaceDir, "memory"), { recursive: true });
44+
await fs.writeFile(
45+
path.join(workspaceDir, "MEMORY.md"),
46+
"Alpha topic about durable index identity\n\nKeep this note for recall.",
47+
);
48+
indexPath = path.join(workspaceDir, "index.sqlite");
49+
});
50+
51+
afterEach(async () => {
52+
if (manager) {
53+
await manager.close();
54+
manager = null;
55+
}
56+
await closeAllMemorySearchManagers();
57+
});
58+
59+
afterAll(async () => {
60+
await closeAllMemorySearchManagers();
61+
if (fixtureRoot) {
62+
await fs.rm(fixtureRoot, { recursive: true, force: true });
63+
}
64+
});
65+
66+
async function createManager(): Promise<MemoryIndexManager> {
67+
const cfg = {
68+
memory: { backend: "builtin" },
69+
agents: {
70+
defaults: {
71+
workspace: workspaceDir,
72+
memorySearch: {
73+
provider: "auto",
74+
model: "",
75+
store: { path: indexPath },
76+
cache: { enabled: false },
77+
sync: { watch: false, onSessionStart: false, onSearch: false },
78+
},
79+
},
80+
list: [{ id: "main", default: true }],
81+
},
82+
} as OpenClawConfig;
83+
const result = await getMemorySearchManager({ cfg, agentId: "main" });
84+
if (!result.manager) {
85+
throw new Error(result.error ?? "manager missing");
86+
}
87+
manager = result.manager as unknown as MemoryIndexManager;
88+
return manager;
89+
}
90+
91+
// Pause a forced full reindex right after `this.db = tempDb` so concurrent
92+
// readers observe the half-built temp DB while the durable index is valid.
93+
function gateNextReindex(memoryManager: MemoryIndexManager): {
94+
entered: Promise<void>;
95+
release: () => void;
96+
} {
97+
const internal = memoryManager as unknown as {
98+
syncMemoryFiles(params: unknown): Promise<void>;
99+
};
100+
const original = internal.syncMemoryFiles.bind(internal);
101+
const enteredGate = deferred();
102+
const releaseGate = deferred();
103+
internal.syncMemoryFiles = async (params: unknown) => {
104+
enteredGate.resolve();
105+
await releaseGate.promise;
106+
internal.syncMemoryFiles = original;
107+
return original(params);
108+
};
109+
return { entered: enteredGate.promise, release: releaseGate.resolve };
110+
}
111+
112+
it("status() reflects the stable durable index during an in-flight full reindex", async () => {
113+
const memoryManager = await createManager();
114+
await memoryManager.sync({ force: true });
115+
116+
const stable = memoryManager.status();
117+
expect(stable.chunks).toBeGreaterThan(0);
118+
119+
const gate = gateNextReindex(memoryManager);
120+
const reindexing = memoryManager.sync({ force: true });
121+
await gate.entered;
122+
123+
// While the temp DB is swapped in mid-build, status() must not report the
124+
// empty/unbuilt index. It should keep reflecting the valid durable index.
125+
const during = memoryManager.status();
126+
gate.release();
127+
await reindexing;
128+
129+
expect(during.chunks).toBe(stable.chunks);
130+
expect(during.files).toBe(stable.files);
131+
});
132+
133+
it("search() returns stable hits instead of reading the half-built temp index", async () => {
134+
const memoryManager = await createManager();
135+
await memoryManager.sync({ force: true });
136+
137+
const baseline = await memoryManager.search("Alpha topic");
138+
expect(baseline.length).toBeGreaterThan(0);
139+
140+
const gate = gateNextReindex(memoryManager);
141+
const reindexing = memoryManager.sync({ force: true });
142+
await gate.entered;
143+
144+
let settled = false;
145+
const searching = memoryManager.search("Alpha topic").then((hits) => {
146+
settled = true;
147+
return hits;
148+
});
149+
await Promise.resolve();
150+
// The search must wait for the reindex rather than read the temp index.
151+
expect(settled).toBe(false);
152+
153+
gate.release();
154+
const hits = await searching;
155+
await reindexing;
156+
expect(hits.length).toBeGreaterThan(0);
157+
});
158+
159+
it("search() waits when provider initialization yields to a full reindex", async () => {
160+
const memoryManager = await createManager();
161+
await memoryManager.sync({ force: true });
162+
163+
const internal = memoryManager as unknown as {
164+
ensureProviderInitialized(): Promise<void>;
165+
};
166+
const originalEnsureProviderInitialized = internal.ensureProviderInitialized.bind(internal);
167+
const gate = gateNextReindex(memoryManager);
168+
let reindexing: Promise<void> = Promise.resolve();
169+
internal.ensureProviderInitialized = async () => {
170+
internal.ensureProviderInitialized = originalEnsureProviderInitialized;
171+
reindexing = memoryManager.sync({ force: true });
172+
await gate.entered;
173+
};
174+
175+
let settled = false;
176+
const searching = memoryManager.search("Alpha topic").then((hits) => {
177+
settled = true;
178+
return hits;
179+
});
180+
await gate.entered;
181+
await Promise.resolve();
182+
183+
expect(settled).toBe(false);
184+
185+
gate.release();
186+
const hits = await searching;
187+
await reindexing;
188+
internal.ensureProviderInitialized = originalEnsureProviderInitialized;
189+
expect(hits.length).toBeGreaterThan(0);
190+
});
191+
192+
it("full reindex waits for active index readers before swapping the DB", async () => {
193+
const memoryManager = await createManager();
194+
await memoryManager.sync({ force: true });
195+
196+
const internal = memoryManager as unknown as {
197+
withIndexRead<T>(read: () => Promise<T>): Promise<T>;
198+
};
199+
const gate = gateNextReindex(memoryManager);
200+
let enteredTempBuild = false;
201+
void gate.entered.then(() => {
202+
enteredTempBuild = true;
203+
});
204+
let reindexing: Promise<void> = Promise.resolve();
205+
206+
await internal.withIndexRead(async () => {
207+
reindexing = memoryManager.sync({ force: true });
208+
await Promise.resolve();
209+
await Promise.resolve();
210+
211+
expect(enteredTempBuild).toBe(false);
212+
});
213+
214+
await gate.entered;
215+
gate.release();
216+
await reindexing;
217+
expect(enteredTempBuild).toBe(true);
218+
});
219+
});

0 commit comments

Comments
 (0)