fix(lbug): run loadCachedEmbeddings reads under withConnLock (#2264 review P2)

loadCachedEmbeddings issued raw conn.query reads on the singleton connection
outside withConnLock — safe today only because it runs before the WAL-checkpoint
driver starts, an ordering invariant not enforced by code. Wrap the whole read in
withConnLock so a future reorder can't race a CHECKPOINT on the connection. Leaf
read; no nested wrapped helpers.

Adds a routing assertion to lbug-conn-serialization.test.ts.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JBJomjoTdBV2eveDVq4JMm
This commit is contained in:
Gergo Magyar 2026-06-21 08:30:57 +00:00
parent 047530d077
commit b96a55fb3c
2 changed files with 70 additions and 54 deletions

View file

@ -1659,67 +1659,75 @@ export const loadCachedEmbeddings = async (): Promise<{
embeddingNodeIds: Set<string>;
embeddings: CachedEmbedding[];
}> => {
if (!conn) {
const c = conn;
if (!c) {
return { embeddingNodeIds: new Set(), embeddings: [] };
}
const embeddingNodeIds = new Set<string>();
const embeddings: CachedEmbedding[] = [];
try {
// Schema migration detection: query with new columns to verify schema version.
// Old schema only had (nodeId, embedding); new schema adds (id, chunkIndex, startLine, endLine, contentHash).
// If the query fails (column missing), we return empty cache to force a full rebuild.
// The whole read runs inside the connection lock (#2264 review P2). It's safe
// today only by call-ordering (loadCachedEmbeddings runs before the WAL driver
// starts), but the lock makes it robust to future reordering — a concurrent
// CHECKPOINT on the singleton connection is the documented corruption trigger.
// Leaf read: no nested withConnLock-wrapped helpers inside.
return withConnLock(async () => {
const embeddingNodeIds = new Set<string>();
const embeddings: CachedEmbedding[] = [];
try {
const check = await conn.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex LIMIT 1`,
);
await readQueryRows(check);
} catch {
return { embeddingNodeIds: new Set(), embeddings: [] };
}
// Try to read contentHash alongside chunk columns
let rows: any;
let hasContentHash = true;
try {
rows = await conn.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding, e.contentHash AS contentHash`,
);
} catch (err: any) {
// Fallback for legacy DBs without contentHash column
const msg = err?.message ?? '';
if (isMissingColumnOrTableError(msg)) {
hasContentHash = false;
rows = await conn.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding`,
// Schema migration detection: query with new columns to verify schema version.
// Old schema only had (nodeId, embedding); new schema adds (id, chunkIndex, startLine, endLine, contentHash).
// If the query fails (column missing), we return empty cache to force a full rebuild.
try {
const check = await c.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex LIMIT 1`,
);
} else {
throw err;
await readQueryRows(check);
} catch {
return { embeddingNodeIds: new Set(), embeddings: [] };
}
}
for (const row of await readQueryRows(rows)) {
const nodeId = String(row.nodeId ?? row[0] ?? '');
if (!nodeId) continue;
embeddingNodeIds.add(nodeId);
const embedding = row.embedding ?? row[4];
if (embedding) {
embeddings.push({
nodeId,
chunkIndex: Number(row.chunkIndex ?? row[1] ?? 0),
startLine: Number(row.startLine ?? row[2] ?? 0),
endLine: Number(row.endLine ?? row[3] ?? 0),
embedding: Array.isArray(embedding)
? embedding.map(Number)
: Array.from(embedding as any).map(Number),
contentHash: hasContentHash ? (row.contentHash ?? row[5] ?? undefined) : undefined,
});
}
}
} catch {
/* embedding table may not exist */
}
return { embeddingNodeIds, embeddings };
// Try to read contentHash alongside chunk columns
let rows: any;
let hasContentHash = true;
try {
rows = await c.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding, e.contentHash AS contentHash`,
);
} catch (err: any) {
// Fallback for legacy DBs without contentHash column
const msg = err?.message ?? '';
if (isMissingColumnOrTableError(msg)) {
hasContentHash = false;
rows = await c.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding`,
);
} else {
throw err;
}
}
for (const row of await readQueryRows(rows)) {
const nodeId = String(row.nodeId ?? row[0] ?? '');
if (!nodeId) continue;
embeddingNodeIds.add(nodeId);
const embedding = row.embedding ?? row[4];
if (embedding) {
embeddings.push({
nodeId,
chunkIndex: Number(row.chunkIndex ?? row[1] ?? 0),
startLine: Number(row.startLine ?? row[2] ?? 0),
endLine: Number(row.endLine ?? row[3] ?? 0),
embedding: Array.isArray(embedding)
? embedding.map(Number)
: Array.from(embedding as any).map(Number),
contentHash: hasContentHash ? (row.contentHash ?? row[5] ?? undefined) : undefined,
});
}
}
} catch {
/* embedding table may not exist */
}
return { embeddingNodeIds, embeddings };
});
};
/**

View file

@ -84,6 +84,14 @@ withTestLbugDB('conn-serialization', () => {
const rows = await adapter.executeQuery('RETURN 1 AS one');
expect(rows).toHaveLength(1);
});
it('U3: loadCachedEmbeddings routes through withConnLock', async () => {
const { loadCachedEmbeddings } = await import('../../src/core/lbug/lbug-adapter.js');
const cached = await loadCachedEmbeddings();
expect(lockSpy).toHaveBeenCalled();
expect(cached.embeddings).toEqual([]);
expect(cached.embeddingNodeIds.size).toBe(0);
});
});
});