Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions extensions/mongoProcessor/MongoQueueProcessor.js
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,7 @@ class MongoQueueProcessor {
},
concurrency: this.mongoProcessorConfig.concurrency,
maxQueued: this.mongoProcessorConfig.maxQueued,
fromOffset: 'earliest',
});
this._consumer.on('error', () => {
MongoProcessorMetrics.onIngestionKafkaConsume('error');
Expand Down
39 changes: 39 additions & 0 deletions tests/unit/mongoProcessor/MongoQueueProcessor.spec.js
Original file line number Diff line number Diff line change
@@ -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);
Expand Down Expand Up @@ -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();
});
});
});
Loading