diff --git a/lib/queuePopulator/IngestionPopulator.js b/lib/queuePopulator/IngestionPopulator.js index 66e9495d9..7aa113a45 100644 --- a/lib/queuePopulator/IngestionPopulator.js +++ b/lib/queuePopulator/IngestionPopulator.js @@ -616,13 +616,31 @@ class IngestionPopulator { const newReaders = this.logReadersUpdate; this.logReadersUpdate = []; async.each(newReaders, (logReader, cb) => logReader.setup(err => { + const zenkoBucket = logReader.getTargetZenkoBucketName(); + // The reader is in neither list while its setup is in flight, so + // its source may have been removed or replaced in the meantime. It + // must then neither be activated nor retried. + if (this._ingestionSources[zenkoBucket] !== logReader) { + this.log.info('skipping log reader after setup: source was removed or replaced', { + method: 'IngestionPopulator._setupUpdatedReaders', + zenkoBucket, + location: logReader.getLocationConstraint(), + }); + return cb(); + } if (err) { // if setup fails for a log reader, don't add it to `logReaders` // log the error and continue setting up others - this.log.fatal('error setting up log reader', { + this.log.error('error setting up log reader, retrying later', { method: 'IngestionPopulator._setupUpdatedReaders', + zenkoBucket, + location: logReader.getLocationConstraint(), error: err, }); + // Setup failures are usually transient (source unreachable, + // invalid credentials...): queue the reader again so that it + // is retried on the next cycle. + this.logReadersUpdate.push(logReader); } else { this.logReaders.push(logReader); } @@ -724,6 +742,12 @@ class IngestionPopulator { * @return {undefined} */ _closeLogState(key) { + // Unregister the source first: a reader whose setup is in flight is in + // neither list, and leaving its entry here would let it register itself + // once the setup completes. + const reader = this._ingestionSources[key]; + delete this._ingestionSources[key]; + if (this._checkAndRemoveLogReader(key, this.logReadersUpdate)) { // if removed from `this.logReadersUpdate`, zookeeper setup has not // been performed yet @@ -731,7 +755,6 @@ class IngestionPopulator { method: 'IngestionPopulator._closeLogState', bucket: key, }); - delete this._ingestionSources[key]; } if (this._checkAndRemoveLogReader(key, this.logReaders)) { // if removed from `this.logReaders`, we must cleanup zookeeper @@ -742,8 +765,6 @@ class IngestionPopulator { }); // we should first validate this reader is not currently processing // a batch of entries before removing zookeeper state - const reader = this._ingestionSources[key]; - delete this._ingestionSources[key]; const path = `${this.ingestionConfig.zookeeperPath}/${key}`; this._removeReaderState(reader, path); } diff --git a/tests/unit/ingestion/IngestionPopulator.js b/tests/unit/ingestion/IngestionPopulator.js index ad10e1904..b48258843 100644 --- a/tests/unit/ingestion/IngestionPopulator.js +++ b/tests/unit/ingestion/IngestionPopulator.js @@ -278,6 +278,207 @@ describe('Ingestion Populator', () => { }); }); + describe('_setupUpdatedReaders', () => { + function addLogReader(zenkoBucket) { + const logReader = sinon.createStubInstance(IngestionReader); + logReader.getTargetZenkoBucketName.returns(zenkoBucket); + ip._ingestionSources[zenkoBucket] = logReader; + ip.logReadersUpdate.push(logReader); + return logReader; + } + + beforeEach(() => { + ip.logReaders = []; + ip.logReadersUpdate = []; + }); + + it('should activate a log reader once its setup succeeds', done => { + const logReader = addLogReader('bucket1'); + logReader.setup.yieldsAsync(null); + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, [logReader]); + assert.deepStrictEqual(ip.logReadersUpdate, []); + done(); + }); + }); + + it('should queue a log reader again when its setup fails', done => { + const logReader = addLogReader('bucket1'); + logReader.setup.yieldsAsync(errors.InternalError); + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, []); + assert.deepStrictEqual(ip.logReadersUpdate, [logReader]); + done(); + }); + }); + + it('should not queue a log reader again when its setup fails and ' + + 'its source is no longer configured', done => { + const logReader = addLogReader('bucket1'); + logReader.setup.yieldsAsync(errors.InternalError); + delete ip._ingestionSources.bucket1; + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, []); + assert.deepStrictEqual(ip.logReadersUpdate, []); + done(); + }); + }); + + it('should not queue a log reader again when its setup fails and ' + + 'its source has been registered with another reader', done => { + const staleReader = addLogReader('bucket1'); + staleReader.setup.yieldsAsync(errors.InternalError); + const currentReader = sinon.createStubInstance(IngestionReader); + currentReader.getTargetZenkoBucketName.returns('bucket1'); + ip._ingestionSources.bucket1 = currentReader; + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, []); + assert.deepStrictEqual(ip.logReadersUpdate, []); + done(); + }); + }); + + it('should not activate a log reader when its setup succeeds and ' + + 'its source is no longer configured', done => { + const logReader = addLogReader('bucket1'); + logReader.setup.yieldsAsync(null); + delete ip._ingestionSources.bucket1; + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, []); + assert.deepStrictEqual(ip.logReadersUpdate, []); + done(); + }); + }); + + it('should not activate a log reader when its setup succeeds and ' + + 'its source has been registered with another reader', done => { + const staleReader = addLogReader('bucket1'); + staleReader.setup.yieldsAsync(null); + const currentReader = sinon.createStubInstance(IngestionReader); + currentReader.getTargetZenkoBucketName.returns('bucket1'); + ip._ingestionSources.bucket1 = currentReader; + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, []); + assert.deepStrictEqual(ip.logReadersUpdate, []); + done(); + }); + }); + + it('should keep setting up other log readers when one fails', done => { + const failingReader = addLogReader('bucket1'); + failingReader.setup.yieldsAsync(errors.InternalError); + const workingReader = addLogReader('bucket2'); + workingReader.setup.yieldsAsync(null); + + ip._setupUpdatedReaders(err => { + assert.ifError(err); + assert.deepStrictEqual(ip.logReaders, [workingReader]); + assert.deepStrictEqual(ip.logReadersUpdate, [failingReader]); + done(); + }); + }); + }); + + describe('_closeLogState', () => { + const ACTIVE_BUCKET = 'active-zenko-bucket'; + const PENDING_BUCKET = 'pending-zenko-bucket'; + + // `IngestionPopulatorMock` stubs out `_closeLogState`, so a real + // populator is needed to exercise it. + let populator; + let removeReaderState; + + beforeEach(() => { + populator = new IngestionPopulator( + null, + zkConfig, + kafkaConfig, + qpConfig, + mConfig, + rConfig, + ingestionConfig, + s3Config, + ); + // `_removeReaderState` polls the reader on a timer that would + // outlive the test + removeReaderState = sinon.stub(populator, '_removeReaderState'); + }); + + /** + * @param {string} zenkoBucket - target zenko bucket of the reader + * @return {object} the stubbed reader + */ + function createReaderMock(zenkoBucket) { + const logReader = sinon.createStubInstance(IngestionReader); + logReader.getTargetZenkoBucketName.returns(zenkoBucket); + return logReader; + } + + it('should unregister a source whose reader is still pending setup', + () => { + const logReader = createReaderMock(PENDING_BUCKET); + populator._ingestionSources[PENDING_BUCKET] = logReader; + populator.logReadersUpdate = [logReader]; + + populator._closeLogState(PENDING_BUCKET); + + assert.deepStrictEqual(populator.logReadersUpdate, []); + assert.strictEqual(populator._ingestionSources[PENDING_BUCKET], + undefined); + // zookeeper state is only created once the setup succeeded + assert.strictEqual(removeReaderState.called, false); + }); + + it('should unregister an active source and clean its zookeeper state', + () => { + const logReader = createReaderMock(ACTIVE_BUCKET); + populator._ingestionSources[ACTIVE_BUCKET] = logReader; + populator.logReaders = [logReader]; + + populator._closeLogState(ACTIVE_BUCKET); + + assert.deepStrictEqual(populator.logReaders, []); + assert.strictEqual(populator._ingestionSources[ACTIVE_BUCKET], + undefined); + // the reader is read before being unregistered, as + // `_removeReaderState` polls it until its batch completes + assert.strictEqual(removeReaderState.calledOnce, true); + assert.strictEqual(removeReaderState.firstCall.args[0], logReader); + assert.strictEqual(removeReaderState.firstCall.args[1], + `${ingestionConfig.zookeeperPath}/${ACTIVE_BUCKET}`); + }); + + it('should unregister a source whose setup is still in flight', () => { + const logReader = createReaderMock(PENDING_BUCKET); + populator._ingestionSources[PENDING_BUCKET] = logReader; + + populator._closeLogState(PENDING_BUCKET); + + assert.strictEqual(populator._ingestionSources[PENDING_BUCKET], + undefined); + assert.strictEqual(removeReaderState.called, false); + }); + + it('should not throw when the source is unknown', () => { + populator._closeLogState('never-configured-bucket'); + + assert.deepStrictEqual(populator._ingestionSources, {}); + assert.strictEqual(removeReaderState.called, false); + }); + }); + describe('_processLogReaderEntries', () => { it('should skip when previous batch currently in progress', () => { const logReaderMock = {