mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-03 02:21:44 +00:00
fix(lbug): harden query result cleanup
This commit is contained in:
parent
694121e9cd
commit
c44356295e
2 changed files with 229 additions and 7 deletions
|
|
@ -218,17 +218,33 @@ 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 Promise.resolve(result.close()).catch(() => {});
|
||||
await closeQueryResult(result);
|
||||
}
|
||||
}
|
||||
if (hasError) throw firstError;
|
||||
};
|
||||
|
||||
const readQueryRows = async (
|
||||
|
|
@ -236,15 +252,23 @@ const readQueryRows = async (
|
|||
): 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 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));
|
||||
}
|
||||
};
|
||||
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue