diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index caa7a58a1..cf8f1cb71 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -218,6 +218,65 @@ const runWithSessionLock = async (operation: () => Promise): Promise => const normalizeCopyPath = (filePath: string): string => filePath.replace(/\\/g, '/'); +const closeQueryResult = async (result: lbug.QueryResult): Promise => { + try { + await result.close(); + } catch { + // Best-effort cleanup only. + } +}; + +const drainQueryResult = async ( + queryResult: lbug.QueryResult | lbug.QueryResult[], +): Promise => { + 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 => { + 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 => { + 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 => { } 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 => { 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 { diff --git a/gitnexus/test/integration/lbug-close-handle-release.test.ts b/gitnexus/test/integration/lbug-close-handle-release.test.ts index c0a3b8758..65a53bc7d 100644 --- a/gitnexus/test/integration/lbug-close-handle-release.test.ts +++ b/gitnexus/test/integration/lbug-close-handle-release.test.ts @@ -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(); + } + }); }); diff --git a/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts b/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts new file mode 100644 index 000000000..3e9ffe8f6 --- /dev/null +++ b/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts @@ -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(); + }); +}); diff --git a/gitnexus/test/unit/lbug-checkpoint.test.ts b/gitnexus/test/unit/lbug-checkpoint.test.ts index 5b68ee9bd..5b9603997 100644 --- a/gitnexus/test/unit/lbug-checkpoint.test.ts +++ b/gitnexus/test/unit/lbug-checkpoint.test.ts @@ -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.