From f42fa75b0f82d88323364d8c4dd89966178da612 Mon Sep 17 00:00:00 2001 From: evolution Date: Tue, 12 May 2026 00:38:53 +0800 Subject: [PATCH] fix(lbug): close query results after reads --- gitnexus/src/core/lbug/lbug-adapter.ts | 93 ++++++++++++------- .../unit/lbug-checkpoint-lifecycle.test.ts | 59 ++++++++++++ 2 files changed, 116 insertions(+), 36 deletions(-) diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 7beaf4293..c011684bc 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -231,6 +231,28 @@ const drainQueryResult = async ( } }; +const readQueryRows = async ( + queryResult: lbug.QueryResult | lbug.QueryResult[], +): Promise => { + 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 => { + 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 => { } 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 { diff --git a/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts b/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts index 1d1362135..e31c2e3ec 100644 --- a/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts +++ b/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts @@ -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(); + }); });