Merge branch 'main' into fix/issue-1486-hook-concurrency-guard

This commit is contained in:
Gergő Magyar 2026-05-12 14:04:14 +01:00 • committed by GitHub
commit 6a34622c38
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 575 additions and 42 deletions

View file

@ -218,6 +218,65 @@ const runWithSessionLock = async <T>(operation: () => Promise<T>): Promise<T> =>
const normalizeCopyPath = (filePath: string): string => filePath.replace(/\\/g, '/');
const closeQueryResult = async (result: lbug.QueryResult): Promise<void> => {
try {
await result.close();
} catch {
// Best-effort cleanup only.
}
};
const drainQueryResult = async (
queryResult: lbug.QueryResult | lbug.QueryResult[],
): Promise<void> => {
const results = Array.isArray(queryResult) ? queryResult : [queryResult];
let firstError: unknown;
let hasError = false;
for (const result of results) {
try {
await result.getAll();
} catch (err) {
if (!hasError) {
firstError = err;
hasError = true;
}
} finally {
await closeQueryResult(result);
}
}
if (hasError) throw firstError;
};
const readQueryRows = async (
queryResult: lbug.QueryResult | lbug.QueryResult[],
): Promise<any[]> => {
const results = Array.isArray(queryResult) ? queryResult : [queryResult];
let rows: any[] = [];
let firstError: unknown;
let hasError = false;
for (let i = 0; i < results.length; i++) {
const result = results[i];
try {
const resultRows = await result.getAll();
if (i === 0) rows = resultRows;
} catch (err) {
if (!hasError) {
firstError = err;
hasError = true;
}
} finally {
await closeQueryResult(result);
}
}
if (hasError) throw firstError;
return rows;
};
const queryAndDrain = async (targetConn: lbug.Connection, cypher: string): Promise<void> => {
const queryResult = await targetConn.query(cypher);
await drainQueryResult(queryResult);
};
export const initLbug = async (dbPath: string) => {
return runWithSessionLock(() => ensureLbugInitialized(dbPath));
};
@ -319,7 +378,7 @@ const doInitLbug = async (dbPath: string) => {
for (const schemaQuery of SCHEMA_QUERIES) {
try {
await conn.query(schemaQuery);
await queryAndDrain(conn, schemaQuery);
} catch (err) {
const msg = err instanceof Error ? err.message : String(err);
// Suppression list:
@ -384,14 +443,14 @@ export const loadGraphToLbug = async (
const copyQuery = getCopyQuery(table, normalizedPath);
try {
await conn.query(copyQuery);
await queryAndDrain(conn, copyQuery);
} catch (err) {
try {
const retryQuery = copyQuery.replace(
'auto_detect=false)',
'auto_detect=false, IGNORE_ERRORS=true)',
);
await conn.query(retryQuery);
await queryAndDrain(conn, retryQuery);
} catch (retryErr) {
const retryMsg = retryErr instanceof Error ? retryErr.message : String(retryErr);
throw new Error(`COPY failed for ${table}: ${retryMsg.slice(0, 200)}`);
@ -433,14 +492,14 @@ export const loadGraphToLbug = async (
}
try {
await conn.query(copyQuery);
await queryAndDrain(conn, copyQuery);
} catch (err) {
try {
const retryQuery = copyQuery.replace(
'auto_detect=false)',
'auto_detect=false, IGNORE_ERRORS=true)',
);
await conn.query(retryQuery);
await queryAndDrain(conn, retryQuery);
} catch (retryErr) {
const retryMsg = retryErr instanceof Error ? retryErr.message : String(retryErr);
warnings.push(`${fromLabel}->${toLabel} (${rows} edges): ${retryMsg.slice(0, 80)}`);
@ -562,11 +621,14 @@ const fallbackRelationshipInserts = async (
const esc = (s: string) =>
s.replace(/'/g, "''").replace(/\\/g, '\\\\').replace(/\n/g, '\\n').replace(/\r/g, '\\r');
await conn.query(`
await queryAndDrain(
conn,
`
MATCH (a:${escapeLabel(fromLabel)} {id: '${esc(fromId)}' }),
(b:${escapeLabel(toLabel)} {id: '${esc(toId)}' })
CREATE (a)-[:${REL_TABLE_NAME} {type: '${esc(relType)}', confidence: ${confidence}, reason: '${esc(reason)}', step: ${step}}]->(b)
`);
`,
);
} catch {
// skip
}
@ -679,14 +741,14 @@ export const insertNodeToLbug = async (
if (targetDbPath) {
const tempHandle = await openLbugConnection(lbug, targetDbPath);
try {
await tempHandle.conn.query(query);
await queryAndDrain(tempHandle.conn, query);
return true;
} finally {
await closeLbugConnection(tempHandle);
}
} else if (conn) {
// Use existing persistent connection (when called from analyze)
await conn.query(query);
await queryAndDrain(conn, query);
return true;
}
@ -757,7 +819,7 @@ export const batchInsertNodesToLbug = async (
query = `MERGE (n:${t} {id: ${escapeValue(properties.id)}}) SET n.name = ${escapeValue(properties.name)}, n.filePath = ${escapeValue(properties.filePath)}, n.startLine = ${properties.startLine || 0}, n.endLine = ${properties.endLine || 0}, n.content = ${escapeValue(properties.content || '')}${descPart}`;
}
await tempConn.query(query);
await queryAndDrain(tempConn, query);
inserted++;
} catch (e: any) {
// Don't console.error here - it corrupts MCP JSON-RPC on stderr
@ -777,11 +839,7 @@ export const executeQuery = async (cypher: string): Promise<any[]> => {
}
const queryResult = await conn.query(cypher);
// LadybugDB uses getAll() instead of hasNext()/getNext()
// Query returns QueryResult for single queries, QueryResult[] for multi-statement
const result = Array.isArray(queryResult) ? queryResult[0] : queryResult;
const rows = await result.getAll();
return rows;
return await readQueryRows(queryResult);
};
export const streamQuery = async (
@ -793,8 +851,10 @@ export const streamQuery = async (
}
const queryResult = await conn.query(cypher);
const result = Array.isArray(queryResult) ? queryResult[0] : queryResult;
const results = Array.isArray(queryResult) ? queryResult : [queryResult];
const result = results[0];
let rowCount = 0;
let streamError: unknown;
try {
while (await result.hasNext()) {
@ -803,11 +863,14 @@ export const streamQuery = async (
rowCount++;
}
return rowCount;
} catch (err) {
streamError = err;
throw err;
} finally {
try {
await result.close();
} catch {
// Best-effort cleanup only.
await drainQueryResult(results);
} catch (err) {
if (streamError === undefined) throw err;
}
}
};
@ -829,8 +892,7 @@ export const executePrepared = async (
throw new Error(`Prepare failed: ${errMsg}`);
}
const queryResult = await conn.execute(stmt, params);
const result = Array.isArray(queryResult) ? queryResult[0] : queryResult;
return await result.getAll();
return await readQueryRows(queryResult);
};
export const executeWithReusedStatement = async (
@ -852,7 +914,7 @@ export const executeWithReusedStatement = async (
}
try {
for (const params of subBatch) {
await conn.execute(stmt, params);
await drainQueryResult(await conn.execute(stmt, params));
}
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
@ -874,8 +936,7 @@ export const getLbugStats = async (): Promise<{ nodes: number; edges: number }>
const queryResult = await conn.query(
`MATCH (n:${escapeTableName(tableName)}) RETURN count(n) AS cnt`,
);
const nodeResult = Array.isArray(queryResult) ? queryResult[0] : queryResult;
const nodeRows = await nodeResult.getAll();
const nodeRows = await readQueryRows(queryResult);
if (nodeRows.length > 0) {
totalNodes += Number(nodeRows[0]?.cnt ?? nodeRows[0]?.[0] ?? 0);
}
@ -889,8 +950,7 @@ export const getLbugStats = async (): Promise<{ nodes: number; edges: number }>
const queryResult = await conn.query(
`MATCH ()-[r:${REL_TABLE_NAME}]->() RETURN count(r) AS cnt`,
);
const edgeResult = Array.isArray(queryResult) ? queryResult[0] : queryResult;
const edgeRows = await edgeResult.getAll();
const edgeRows = await readQueryRows(queryResult);
if (edgeRows.length > 0) {
totalEdges = Number(edgeRows[0]?.cnt ?? edgeRows[0]?.[0] ?? 0);
}
@ -926,8 +986,7 @@ export const loadCachedEmbeddings = async (): Promise<{
const check = await conn.query(
`MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex LIMIT 1`,
);
const checkResult = Array.isArray(check) ? check[0] : check;
await checkResult.getAll();
await readQueryRows(check);
} catch {
return { embeddingNodeIds: new Set(), embeddings: [] };
}
@ -951,8 +1010,7 @@ export const loadCachedEmbeddings = async (): Promise<{
throw err;
}
}
const result = Array.isArray(rows) ? rows[0] : rows;
for (const row of await result.getAll()) {
for (const row of await readQueryRows(rows)) {
const nodeId = String(row.nodeId ?? row[0] ?? '');
if (!nodeId) continue;
embeddingNodeIds.add(nodeId);
@ -1060,7 +1118,8 @@ export const fetchExistingEmbeddingHashes = async (
export const flushWAL = async (): Promise<void> => {
if (!conn) return;
try {
await conn.query('CHECKPOINT');
const checkpointResult = await conn.query('CHECKPOINT');
await drainQueryResult(checkpointResult);
} catch {
/* ignore — older LadybugDB or schemaless DB may not accept it */
}
@ -1170,13 +1229,13 @@ export const deleteNodesForFile = async (
const countResult = await targetConn!.query(
`MATCH (n:${tn}) WHERE n.filePath = '${escapedPath}' RETURN count(n) AS cnt`,
);
const result = Array.isArray(countResult) ? countResult[0] : countResult;
const rows = await result.getAll();
const rows = await readQueryRows(countResult);
const count = Number(rows[0]?.cnt ?? rows[0]?.[0] ?? 0);
if (count > 0) {
// Delete nodes (and implicitly their relationships via DETACH)
await targetConn!.query(
await queryAndDrain(
targetConn!,
`MATCH (n:${tn}) WHERE n.filePath = '${escapedPath}' DETACH DELETE n`,
);
deletedNodes += count;
@ -1188,7 +1247,8 @@ export const deleteNodesForFile = async (
// Also delete any embeddings for nodes in this file
try {
await targetConn!.query(
await queryAndDrain(
targetConn!,
`MATCH (e:${EMBEDDING_TABLE_NAME}) WHERE e.nodeId STARTS WITH '${escapedPath}' DELETE e`,
);
} catch {
@ -1303,7 +1363,7 @@ export const loadFTSExtension = async (
throw new Error('LadybugDB not initialized. Call initLbug first.');
}
const loaded = await extensionManager.ensure((sql) => c.query(sql), 'fts', 'FTS', opts);
const loaded = await extensionManager.ensure((sql) => queryAndDrain(c, sql), 'fts', 'FTS', opts);
if (loaded && useModuleState) ftsLoaded = true;
return loaded;
};
@ -1333,7 +1393,12 @@ export const loadVectorExtension = async (
throw new Error('LadybugDB not initialized. Call initLbug first.');
}
const loaded = await extensionManager.ensure((sql) => c.query(sql), 'VECTOR', 'VECTOR', opts);
const loaded = await extensionManager.ensure(
(sql) => queryAndDrain(c, sql),
'VECTOR',
'VECTOR',
opts,
);
if (loaded && useModuleState) vectorExtensionLoaded = true;
return loaded;
};
@ -1365,7 +1430,7 @@ export const createFTSIndex = async (
const query = `CALL CREATE_FTS_INDEX('${tableName}', '${indexName}', [${propList}], stemmer := '${stemmer}')`;
try {
await conn.query(query);
await queryAndDrain(conn, query);
ensuredFTSIndexes.add(key);
} catch (e: any) {
if (e.message?.includes('already exists')) {
@ -1449,8 +1514,7 @@ export const queryFTS = async (
try {
const queryResult = await conn.query(cypher);
const result = Array.isArray(queryResult) ? queryResult[0] : queryResult;
const rows = await result.getAll();
const rows = await readQueryRows(queryResult);
return rows.map((row: any) => {
const node = row.node || row[0] || {};
@ -1481,7 +1545,7 @@ export const dropFTSIndex = async (tableName: string, indexName: string): Promis
}
try {
await conn.query(`CALL DROP_FTS_INDEX('${tableName}', '${indexName}')`);
await queryAndDrain(conn, `CALL DROP_FTS_INDEX('${tableName}', '${indexName}')`);
} catch {
// Index may not exist
} finally {

View file

@ -8,9 +8,16 @@
* absorbed by the open-time retry in `lbug-config.ts`.
*/
import path from 'path';
import { describe, it } from 'vitest';
import { describe, expect, it } from 'vitest';
import { createTempDir } from '../helpers/test-db.js';
/**
* LadybugDB's native Windows file lock can outlive Database.close() for
* same-process close/reopen cycles. Keep true reopen coverage on POSIX and
* cover ordering deterministically in lbug-checkpoint-lifecycle.test.ts.
*/
const itLbugReopen = process.platform === 'win32' ? it.skip : it;
describe('safeClose — close + reopen does not surface lock errors', () => {
it('survives 10 sequential open/close/reopen cycles on the same path', async () => {
const tmp = await createTempDir('gitnexus-lbug-close-cycle-');
@ -38,4 +45,38 @@ describe('safeClose — close + reopen does not surface lock errors', () => {
await tmp.cleanup();
}
});
itLbugReopen('flushes WAL when switching between two database paths in one process', async () => {
const repoA = await createTempDir('gitnexus-lbug-switch-a-');
const repoB = await createTempDir('gitnexus-lbug-switch-b-');
const dbPathA = path.join(repoA.dbPath, 'lbug');
const dbPathB = path.join(repoB.dbPath, 'lbug');
try {
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.withLbugDb(dbPathA, async () => {
await adapter.executeQuery(
"CREATE (:File {id: 'file:a', name: 'a.ts', filePath: 'a.ts', content: 'repo a'})",
);
});
await adapter.withLbugDb(dbPathB, async () => {
await adapter.executeQuery(
"CREATE (:File {id: 'file:b', name: 'b.ts', filePath: 'b.ts', content: 'repo b'})",
);
});
const rows = await adapter.withLbugDb(dbPathA, async () =>
adapter.executeQuery("MATCH (n:File {id: 'file:a'}) RETURN n.filePath AS filePath"),
);
expect(rows).toEqual([{ filePath: 'a.ts' }]);
} finally {
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.closeLbug().catch(() => {});
await repoA.cleanup();
await repoB.cleanup();
}
});
});

View file

@ -0,0 +1,420 @@
import { afterEach, describe, expect, it, vi } from 'vitest';
describe('lbug adapter CHECKPOINT lifecycle', () => {
afterEach(() => {
vi.doUnmock('../../src/core/lbug/lbug-config.js');
vi.doUnmock('../../src/core/lbug/extension-loader.js');
vi.resetModules();
vi.clearAllMocks();
});
it('drains and closes CHECKPOINT result before closing connection and database handles', async () => {
vi.resetModules();
const events: string[] = [];
const checkpointResult = {
getAll: vi.fn(async () => {
events.push('checkpoint:getAll');
return [];
}),
close: vi.fn(() => {
events.push('checkpoint:close');
}),
};
const genericResult = {
getAll: vi.fn(async () => []),
close: vi.fn(),
};
const conn = {
query: vi.fn(async (sql: string) => {
if (sql === 'CHECKPOINT') {
events.push('checkpoint:query');
return checkpointResult;
}
return genericResult;
}),
close: vi.fn(async () => {
events.push('conn:close');
}),
};
const db = {
close: vi.fn(async () => {
events.push('db:close');
}),
};
vi.doMock('../../src/core/lbug/lbug-config.js', () => ({
openLbugConnection: vi.fn(async () => ({ db, conn })),
closeLbugConnection: vi.fn(async () => {}),
isDbBusyError: vi.fn((err: unknown) => String(err).toLowerCase().includes('lock')),
isOpenRetryExhausted: vi.fn(() => false),
waitForWindowsHandleRelease: vi.fn(async () => true),
}));
vi.doMock('../../src/core/lbug/extension-loader.js', () => ({
extensionManager: {
ensure: vi.fn(async () => true),
getCapabilities: vi.fn(() => []),
reset: vi.fn(),
},
}));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.initLbug('/tmp/gitnexus-lbug-checkpoint-lifecycle/lbug');
events.length = 0;
await adapter.closeLbug();
expect(events).toEqual([
'checkpoint:query',
'checkpoint:getAll',
'checkpoint:close',
'conn:close',
'db:close',
]);
});
it('closes normal query results after reading rows', async () => {
vi.resetModules();
const events: string[] = [];
const queryResult = {
getAll: vi.fn(async () => {
events.push('query:getAll');
return [{ id: 'file:a' }];
}),
close: vi.fn(() => {
events.push('query:close');
}),
};
const genericResult = {
getAll: vi.fn(async () => []),
close: vi.fn(),
};
const conn = {
query: vi.fn(async (sql: string) => {
if (sql === 'MATCH (n:File) RETURN n.id AS id') {
events.push('query:run');
return queryResult;
}
return genericResult;
}),
close: vi.fn(async () => {}),
};
const db = {
close: vi.fn(async () => {}),
};
vi.doMock('../../src/core/lbug/lbug-config.js', () => ({
openLbugConnection: vi.fn(async () => ({ db, conn })),
closeLbugConnection: vi.fn(async () => {}),
isDbBusyError: vi.fn((err: unknown) => String(err).toLowerCase().includes('lock')),
isOpenRetryExhausted: vi.fn(() => false),
waitForWindowsHandleRelease: vi.fn(async () => true),
}));
vi.doMock('../../src/core/lbug/extension-loader.js', () => ({
extensionManager: {
ensure: vi.fn(async () => true),
getCapabilities: vi.fn(() => []),
reset: vi.fn(),
},
}));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.initLbug('/tmp/gitnexus-lbug-query-lifecycle/lbug');
events.length = 0;
await expect(adapter.executeQuery('MATCH (n:File) RETURN n.id AS id')).resolves.toEqual([
{ id: 'file:a' },
]);
expect(events).toEqual(['query:run', 'query:getAll', 'query:close']);
await adapter.closeLbug();
});
it('treats synchronous query result close errors as best-effort cleanup', async () => {
vi.resetModules();
const queryResult = {
getAll: vi.fn(async () => [{ id: 'file:a' }]),
close: vi.fn(() => {
throw new Error('close failed');
}),
};
const genericResult = {
getAll: vi.fn(async () => []),
close: vi.fn(),
};
const conn = {
query: vi.fn(async (sql: string) => {
if (sql === 'MATCH (n:File) RETURN n.id AS id') {
return queryResult;
}
return genericResult;
}),
close: vi.fn(async () => {}),
};
const db = {
close: vi.fn(async () => {}),
};
vi.doMock('../../src/core/lbug/lbug-config.js', () => ({
openLbugConnection: vi.fn(async () => ({ db, conn })),
closeLbugConnection: vi.fn(async () => {}),
isDbBusyError: vi.fn((err: unknown) => String(err).toLowerCase().includes('lock')),
isOpenRetryExhausted: vi.fn(() => false),
waitForWindowsHandleRelease: vi.fn(async () => true),
}));
vi.doMock('../../src/core/lbug/extension-loader.js', () => ({
extensionManager: {
ensure: vi.fn(async () => true),
getCapabilities: vi.fn(() => []),
reset: vi.fn(),
},
}));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.initLbug('/tmp/gitnexus-lbug-sync-close-lifecycle/lbug');
await expect(adapter.executeQuery('MATCH (n:File) RETURN n.id AS id')).resolves.toEqual([
{ id: 'file:a' },
]);
expect(queryResult.close).toHaveBeenCalledOnce();
await adapter.closeLbug();
});
it('closes later query results when an earlier array result fails to read', async () => {
vi.resetModules();
const events: string[] = [];
const firstResult = {
getAll: vi.fn(async () => {
events.push('first:getAll');
throw new Error('read failed');
}),
close: vi.fn(() => {
events.push('first:close');
}),
};
const secondResult = {
getAll: vi.fn(async () => {
events.push('second:getAll');
return [];
}),
close: vi.fn(() => {
events.push('second:close');
}),
};
const genericResult = {
getAll: vi.fn(async () => []),
close: vi.fn(),
};
const conn = {
query: vi.fn(async (sql: string) => {
if (sql === 'MATCH (n:File) RETURN n.id AS id') {
return [firstResult, secondResult];
}
return genericResult;
}),
close: vi.fn(async () => {}),
};
const db = {
close: vi.fn(async () => {}),
};
vi.doMock('../../src/core/lbug/lbug-config.js', () => ({
openLbugConnection: vi.fn(async () => ({ db, conn })),
closeLbugConnection: vi.fn(async () => {}),
isDbBusyError: vi.fn((err: unknown) => String(err).toLowerCase().includes('lock')),
isOpenRetryExhausted: vi.fn(() => false),
waitForWindowsHandleRelease: vi.fn(async () => true),
}));
vi.doMock('../../src/core/lbug/extension-loader.js', () => ({
extensionManager: {
ensure: vi.fn(async () => true),
getCapabilities: vi.fn(() => []),
reset: vi.fn(),
},
}));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.initLbug('/tmp/gitnexus-lbug-array-error-lifecycle/lbug');
await expect(adapter.executeQuery('MATCH (n:File) RETURN n.id AS id')).rejects.toThrow(
'read failed',
);
expect(events).toEqual(['first:getAll', 'first:close', 'second:getAll', 'second:close']);
await adapter.closeLbug();
});
it('closes non-first stream query results when LadybugDB returns an array', async () => {
vi.resetModules();
const events: string[] = [];
const firstResult = {
hasNext: vi
.fn()
.mockImplementationOnce(() => {
events.push('first:hasNext:true');
return true;
})
.mockImplementationOnce(() => {
events.push('first:hasNext:false');
return false;
}),
getNext: vi.fn(async () => {
events.push('first:getNext');
return { id: 'file:a' };
}),
getAll: vi.fn(async () => {
events.push('first:getAll');
return [];
}),
close: vi.fn(() => {
events.push('first:close');
}),
};
const secondResult = {
getAll: vi.fn(async () => {
events.push('second:getAll');
return [];
}),
close: vi.fn(() => {
events.push('second:close');
}),
};
const genericResult = {
getAll: vi.fn(async () => []),
close: vi.fn(),
};
const conn = {
query: vi.fn(async (sql: string) => {
if (sql === 'MATCH (n:File) RETURN n.id AS id') {
events.push('stream:query');
return [firstResult, secondResult];
}
return genericResult;
}),
close: vi.fn(async () => {}),
};
const db = {
close: vi.fn(async () => {}),
};
vi.doMock('../../src/core/lbug/lbug-config.js', () => ({
openLbugConnection: vi.fn(async () => ({ db, conn })),
closeLbugConnection: vi.fn(async () => {}),
isDbBusyError: vi.fn((err: unknown) => String(err).toLowerCase().includes('lock')),
isOpenRetryExhausted: vi.fn(() => false),
waitForWindowsHandleRelease: vi.fn(async () => true),
}));
vi.doMock('../../src/core/lbug/extension-loader.js', () => ({
extensionManager: {
ensure: vi.fn(async () => true),
getCapabilities: vi.fn(() => []),
reset: vi.fn(),
},
}));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.initLbug('/tmp/gitnexus-lbug-stream-lifecycle/lbug');
const rows: unknown[] = [];
events.length = 0;
await expect(
adapter.streamQuery('MATCH (n:File) RETURN n.id AS id', (row) => {
rows.push(row);
}),
).resolves.toBe(1);
expect(rows).toEqual([{ id: 'file:a' }]);
expect(events).toEqual([
'stream:query',
'first:hasNext:true',
'first:getNext',
'first:hasNext:false',
'first:getAll',
'first:close',
'second:getAll',
'second:close',
]);
await adapter.closeLbug();
});
it('drains stream query results when row handling fails before the result is exhausted', async () => {
vi.resetModules();
const events: string[] = [];
const queryResult = {
hasNext: vi.fn(() => {
events.push('stream:hasNext');
return true;
}),
getNext: vi.fn(async () => {
events.push('stream:getNext');
return { id: 'file:a' };
}),
getAll: vi.fn(async () => {
events.push('stream:getAll');
return [{ id: 'file:b' }];
}),
close: vi.fn(() => {
events.push('stream:close');
}),
};
const genericResult = {
getAll: vi.fn(async () => []),
close: vi.fn(),
};
const conn = {
query: vi.fn(async (sql: string) => {
if (sql === 'MATCH (n:File) RETURN n.id AS id') {
events.push('stream:query');
return queryResult;
}
return genericResult;
}),
close: vi.fn(async () => {}),
};
const db = {
close: vi.fn(async () => {}),
};
vi.doMock('../../src/core/lbug/lbug-config.js', () => ({
openLbugConnection: vi.fn(async () => ({ db, conn })),
closeLbugConnection: vi.fn(async () => {}),
isDbBusyError: vi.fn((err: unknown) => String(err).toLowerCase().includes('lock')),
isOpenRetryExhausted: vi.fn(() => false),
waitForWindowsHandleRelease: vi.fn(async () => true),
}));
vi.doMock('../../src/core/lbug/extension-loader.js', () => ({
extensionManager: {
ensure: vi.fn(async () => true),
getCapabilities: vi.fn(() => []),
reset: vi.fn(),
},
}));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
await adapter.initLbug('/tmp/gitnexus-lbug-stream-error-lifecycle/lbug');
await expect(
adapter.streamQuery('MATCH (n:File) RETURN n.id AS id', () => {
throw new Error('client disconnected');
}),
).rejects.toThrow('client disconnected');
expect(events).toEqual([
'stream:query',
'stream:hasNext',
'stream:getNext',
'stream:getAll',
'stream:close',
]);
await adapter.closeLbug();
});
});

View file

@ -58,6 +58,14 @@ describe('flushWAL / safeClose — consolidation guard (#1376)', () => {
expect(matches.length).toBe(1);
});
it('flushWAL drains and closes the CHECKPOINT result before returning', () => {
const flushBody = adapterSource.slice(
adapterSource.indexOf('export const flushWAL'),
adapterSource.indexOf('export const safeClose'),
);
expect(flushBody).toMatch(/await drainQueryResult\(checkpointResult\)/);
});
it('conn.close() only appears inside safeClose (with eslint-disable)', () => {
// Every conn.close() in the adapter must live inside safeClose, guarded
// by the eslint-disable comment. Count occurrences to catch leaks.