Merge pull request #17292 from iamEAP/search/close-pg-connections-on-error
Fix PG search indexer connection exhaustion issue
This commit is contained in:
@@ -0,0 +1,7 @@
|
||||
---
|
||||
'@backstage/plugin-search-backend-module-pg': patch
|
||||
---
|
||||
|
||||
Fixed a bug that could cause orphaned PG connections to accumulate (eventually
|
||||
exhausting available connections) when errors were encountered earlier in the
|
||||
search indexing process.
|
||||
@@ -285,6 +285,9 @@ export class ElasticSearchSearchEngine implements SearchEngine {
|
||||
});
|
||||
|
||||
// Attempt cleanup upon failure.
|
||||
// todo(@backstage/discoverability-maintainers): Consider introducing a more
|
||||
// formal mechanism for handling such errors in BatchSearchEngineIndexer and
|
||||
// replacing this handler with it. See: #17291
|
||||
indexer.on('error', async e => {
|
||||
indexerLogger.error(`Failed to index documents for type ${type}`, e);
|
||||
let cleanupError: Error | undefined;
|
||||
|
||||
@@ -15,6 +15,7 @@
|
||||
*/
|
||||
import { TestPipeline } from '@backstage/plugin-search-backend-node';
|
||||
import { range } from 'lodash';
|
||||
import { Transform } from 'stream';
|
||||
import { PgSearchEngineIndexer } from './PgSearchEngineIndexer';
|
||||
import { DatabaseStore } from '../database';
|
||||
|
||||
@@ -22,6 +23,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
const tx = {
|
||||
rollback: jest.fn(),
|
||||
commit: jest.fn(),
|
||||
isCompleted: jest.fn(),
|
||||
} as any;
|
||||
let database: jest.Mocked<DatabaseStore>;
|
||||
let indexer: PgSearchEngineIndexer;
|
||||
@@ -36,6 +38,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
completeInsert: jest.fn(),
|
||||
prepareInsert: jest.fn(),
|
||||
};
|
||||
tx.isCompleted.mockReturnValue(false);
|
||||
indexer = new PgSearchEngineIndexer({
|
||||
batchSize: 100,
|
||||
type: 'my-type',
|
||||
@@ -64,6 +67,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
);
|
||||
expect(database.completeInsert).toHaveBeenCalledWith(tx, 'my-type');
|
||||
expect(tx.commit).toHaveBeenCalled();
|
||||
expect(tx.rollback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('should batch insert documents', async () => {
|
||||
@@ -87,7 +91,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
expect(database.getTransaction).toHaveBeenCalledTimes(1);
|
||||
expect(database.insertDocuments).not.toHaveBeenCalled();
|
||||
expect(database.completeInsert).not.toHaveBeenCalled();
|
||||
expect(tx.rollback).toHaveBeenCalled();
|
||||
expect(tx.rollback).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('should close out stream and bubble up error on prepare', async () => {
|
||||
@@ -100,6 +104,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
},
|
||||
];
|
||||
|
||||
tx.isCompleted.mockReturnValue(true);
|
||||
database.prepareInsert.mockRejectedValueOnce(expectedError);
|
||||
const result = await TestPipeline.fromIndexer(indexer)
|
||||
.withDocuments(documents)
|
||||
@@ -109,6 +114,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
expect(database.insertDocuments).not.toHaveBeenCalled();
|
||||
expect(database.completeInsert).not.toHaveBeenCalled();
|
||||
expect(result.error).toBe(expectedError);
|
||||
expect(tx.rollback).toHaveBeenCalledTimes(1);
|
||||
expect(tx.rollback).toHaveBeenCalledWith(expectedError);
|
||||
});
|
||||
|
||||
@@ -122,6 +128,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
},
|
||||
];
|
||||
|
||||
tx.isCompleted.mockReturnValue(true);
|
||||
database.insertDocuments.mockRejectedValueOnce(expectedError);
|
||||
const result = await TestPipeline.fromIndexer(indexer)
|
||||
.withDocuments(documents)
|
||||
@@ -131,6 +138,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
expect(database.prepareInsert).toHaveBeenCalledTimes(1);
|
||||
expect(database.completeInsert).not.toHaveBeenCalled();
|
||||
expect(result.error).toBe(expectedError);
|
||||
expect(tx.rollback).toHaveBeenCalledTimes(1);
|
||||
expect(tx.rollback).toHaveBeenCalledWith(expectedError);
|
||||
});
|
||||
|
||||
@@ -144,6 +152,7 @@ describe('PgSearchEngineIndexer', () => {
|
||||
},
|
||||
];
|
||||
|
||||
tx.isCompleted.mockReturnValue(true);
|
||||
database.completeInsert.mockRejectedValueOnce(expectedError);
|
||||
const result = await TestPipeline.fromIndexer(indexer)
|
||||
.withDocuments(documents)
|
||||
@@ -154,6 +163,41 @@ describe('PgSearchEngineIndexer', () => {
|
||||
expect(database.insertDocuments).toHaveBeenCalledTimes(1);
|
||||
expect(database.completeInsert).toHaveBeenCalledTimes(1);
|
||||
expect(result.error).toBe(expectedError);
|
||||
expect(tx.rollback).toHaveBeenCalledTimes(1);
|
||||
expect(tx.rollback).toHaveBeenCalledWith(expectedError);
|
||||
});
|
||||
|
||||
it('should rollback transaction on upstream error', async () => {
|
||||
// Given a decorator that results in an error
|
||||
let counter = 0;
|
||||
const expectedError = new Error('Upstream error');
|
||||
const errorDecorator = new Transform({ objectMode: true });
|
||||
errorDecorator._transform = (chunk, _enc, cb) => {
|
||||
counter++;
|
||||
if (counter > 1) {
|
||||
cb(expectedError);
|
||||
} else {
|
||||
cb(undefined, chunk);
|
||||
}
|
||||
};
|
||||
|
||||
// When the decorator is run in a pipeline with the PG indexer
|
||||
const result = await TestPipeline.fromIndexer(indexer)
|
||||
.withDecorator(errorDecorator)
|
||||
.withDocuments([
|
||||
{ title: 'a', text: 'a', location: '/a' },
|
||||
{ title: 'b', text: 'b', location: '/b' },
|
||||
])
|
||||
.execute();
|
||||
|
||||
// And we allow async teardown logic to complete
|
||||
await new Promise(resolve => setImmediate(resolve));
|
||||
|
||||
// Then the transaction should have been closed with the expected error.
|
||||
expect(database.getTransaction).toHaveBeenCalledTimes(1);
|
||||
expect(database.completeInsert).not.toHaveBeenCalled();
|
||||
expect(result.error).toBe(expectedError);
|
||||
expect(tx.rollback).toHaveBeenCalledTimes(1);
|
||||
expect(tx.rollback).toHaveBeenCalledWith(expectedError);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -92,4 +92,31 @@ export class PgSearchEngineIndexer extends BatchSearchEngineIndexer {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Custom handler covering the case where an error occurred somewhere else in
|
||||
* the indexing pipeline (e.g. a collator or decorator). In such cases, the
|
||||
* finalize method is not called, which leaves a dangling transaction and
|
||||
* therefore an open connection to PG. This handler ensures we close the
|
||||
* transaction and associated connection.
|
||||
*
|
||||
* todo(@backstage/discoverability-maintainers): Consider introducing a more
|
||||
* formal mechanism for handling such errors in BatchSearchEngineIndexer and
|
||||
* replacing this method with it. See: #17291
|
||||
*
|
||||
* @internal
|
||||
*/
|
||||
async _destroy(error: Error | null, done: (error?: Error | null) => void) {
|
||||
// Ignore situations where there was no error.
|
||||
if (!error) {
|
||||
done();
|
||||
return;
|
||||
}
|
||||
|
||||
if (!this.tx!.isCompleted()) {
|
||||
await this.tx!.rollback(error);
|
||||
}
|
||||
|
||||
done(error);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user