diff --git a/extensions/mongoProcessor/MongoQueueProcessor.js b/extensions/mongoProcessor/MongoQueueProcessor.js index cc2843f07..09b465832 100644 --- a/extensions/mongoProcessor/MongoQueueProcessor.js +++ b/extensions/mongoProcessor/MongoQueueProcessor.js @@ -147,6 +147,7 @@ class MongoQueueProcessor { }, concurrency: this.mongoProcessorConfig.concurrency, maxQueued: this.mongoProcessorConfig.maxQueued, + fromOffset: 'earliest', }); this._consumer.on('error', () => { MongoProcessorMetrics.onIngestionKafkaConsume('error'); diff --git a/tests/unit/mongoProcessor/MongoQueueProcessor.spec.js b/tests/unit/mongoProcessor/MongoQueueProcessor.spec.js index 4ed3b85a1..b58f77a79 100644 --- a/tests/unit/mongoProcessor/MongoQueueProcessor.spec.js +++ b/tests/unit/mongoProcessor/MongoQueueProcessor.spec.js @@ -1,9 +1,12 @@ const assert = require('assert'); +const sinon = require('sinon'); const MongoQueueProcessor = require('../../../extensions/mongoProcessor/MongoQueueProcessor'); const ObjectQueueEntry = require('../../../lib/models/ObjectQueueEntry'); +const BackbeatConsumer = require('../../../lib/BackbeatConsumer'); +const Config = require('../../../lib/Config'); function _makeProcessor(bootstrapList) { const proc = Object.create(MongoQueueProcessor.prototype); @@ -234,3 +237,39 @@ describe('MongoQueueProcessor._updateReplicationInfo', () => { assert.strictEqual(bySite['cloud-b'].status, 'PENDING'); }); }); + +describe('MongoQueueProcessor.start', () => { + afterEach(() => { + sinon.restore(); + }); + + function startProcessor() { + sinon.stub(BackbeatConsumer.prototype, '_init'); + sinon.stub(Config, 'getBootstrapList').returns([]); + sinon.stub(Config, 'on'); + + const proc = _makeProcessor([]); + proc.logger = { info: () => {}, error: () => {}, fatal: () => {} }; + proc.kafkaConfig = { hosts: 'localhost:9092' }; + proc.mongoProcessorConfig = { + topic: 'backbeat-ingestion', + groupId: 'backbeat-ingestion-group', + concurrency: 5, + }; + proc._setupMetricsClients = cb => cb(); + proc._mongoClient = { setup: cb => cb() }; + + proc.start(); + + return proc; + } + + it('starts the consumer at the earliest offset', done => { + const proc = startProcessor(); + + setImmediate(() => { + assert.strictEqual(proc._consumer._fromOffset, 'earliest'); + done(); + }); + }); +});