From c44356295ecfde5848b918ff9ee52c74334737a2 Mon Sep 17 00:00:00 2001 From: evolution Date: Tue, 12 May 2026 11:43:54 +0800 Subject: [PATCH] fix(lbug): harden query result cleanup --- gitnexus/src/core/lbug/lbug-adapter.ts | 39 +++- .../unit/lbug-checkpoint-lifecycle.test.ts | 197 +++++++++++++++++- 2 files changed, 229 insertions(+), 7 deletions(-) diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 2405d677f..e7043b7f2 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -218,17 +218,33 @@ 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 Promise.resolve(result.close()).catch(() => {}); + await closeQueryResult(result); } } + if (hasError) throw firstError; }; const readQueryRows = async ( @@ -236,15 +252,23 @@ const readQueryRows = async ( ): 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 Promise.resolve(result.close()).catch(() => {}); + await closeQueryResult(result); } } + if (hasError) throw firstError; return rows; }; @@ -830,6 +854,7 @@ export const streamQuery = async ( const results = Array.isArray(queryResult) ? queryResult : [queryResult]; const result = results[0]; let rowCount = 0; + let streamError: unknown; try { while (await result.hasNext()) { @@ -838,13 +863,15 @@ 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; } - await drainQueryResult(results.slice(1)); } }; diff --git a/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts b/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts index 4b2dd82ad..3e9ffe8f6 100644 --- a/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts +++ b/gitnexus/test/unit/lbug-checkpoint-lifecycle.test.ts @@ -132,6 +132,123 @@ describe('lbug adapter CHECKPOINT lifecycle', () => { 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(); @@ -151,7 +268,10 @@ describe('lbug adapter CHECKPOINT lifecycle', () => { events.push('first:getNext'); return { id: 'file:a' }; }), - getAll: vi.fn(async () => []), + getAll: vi.fn(async () => { + events.push('first:getAll'); + return []; + }), close: vi.fn(() => { events.push('first:close'); }), @@ -215,6 +335,7 @@ describe('lbug adapter CHECKPOINT lifecycle', () => { 'first:hasNext:true', 'first:getNext', 'first:hasNext:false', + 'first:getAll', 'first:close', 'second:getAll', 'second:close', @@ -222,4 +343,78 @@ describe('lbug adapter CHECKPOINT lifecycle', () => { 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(); + }); });