mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-03 02:21:44 +00:00
fix(lbug): close query results after reads
This commit is contained in:
parent
924554e36f
commit
f42fa75b0f
2 changed files with 116 additions and 36 deletions
|
|
@ -231,6 +231,28 @@ const drainQueryResult = async (
|
|||
}
|
||||
};
|
||||
|
||||
const readQueryRows = async (
|
||||
queryResult: lbug.QueryResult | lbug.QueryResult[],
|
||||
): Promise<any[]> => {
|
||||
const results = Array.isArray(queryResult) ? queryResult : [queryResult];
|
||||
let rows: any[] = [];
|
||||
for (let i = 0; i < results.length; i++) {
|
||||
const result = results[i];
|
||||
try {
|
||||
const resultRows = await result.getAll();
|
||||
if (i === 0) rows = resultRows;
|
||||
} finally {
|
||||
await Promise.resolve(result.close()).catch(() => {});
|
||||
}
|
||||
}
|
||||
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));
|
||||
};
|
||||
|
|
@ -332,7 +354,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:
|
||||
|
|
@ -397,14 +419,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)}`);
|
||||
|
|
@ -446,14 +468,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)}`);
|
||||
|
|
@ -575,11 +597,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
|
||||
}
|
||||
|
|
@ -692,14 +717,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;
|
||||
}
|
||||
|
||||
|
|
@ -770,7 +795,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
|
||||
|
|
@ -790,11 +815,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 (
|
||||
|
|
@ -842,8 +863,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 (
|
||||
|
|
@ -865,7 +885,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);
|
||||
|
|
@ -887,8 +907,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);
|
||||
}
|
||||
|
|
@ -902,8 +921,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);
|
||||
}
|
||||
|
|
@ -939,8 +957,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: [] };
|
||||
}
|
||||
|
|
@ -964,8 +981,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);
|
||||
|
|
@ -1184,13 +1200,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;
|
||||
|
|
@ -1202,7 +1218,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 {
|
||||
|
|
@ -1246,7 +1263,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;
|
||||
};
|
||||
|
|
@ -1276,7 +1293,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;
|
||||
};
|
||||
|
|
@ -1308,7 +1330,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')) {
|
||||
|
|
@ -1392,8 +1414,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] || {};
|
||||
|
|
@ -1424,7 +1445,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 {
|
||||
|
|
|
|||
|
|
@ -72,4 +72,63 @@ describe('lbug adapter CHECKPOINT lifecycle', () => {
|
|||
'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();
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue