diff --git a/lib/models/QueueEntry.js b/lib/models/QueueEntry.js index 516416e3b..3d0049f51 100644 --- a/lib/models/QueueEntry.js +++ b/lib/models/QueueEntry.js @@ -5,6 +5,26 @@ const BucketQueueEntry = require('./BucketQueueEntry'); const BucketMdQueueEntry = require('./BucketMdQueueEntry'); const DeleteOpQueueEntry = require('./DeleteOpQueueEntry'); +/** + * Decode a nested field of a kafka entry, which producers may send either + * as a JSON string or as an already-decoded object. + * + * @param {string|object} field - JSON string or decoded object + * @return {object} the decoded field + * @throws {Error} if the field is neither a JSON string nor an object; the + * caller turns this into a "malformed JSON in kafka entry" error + */ +function decodeField(field) { + if (typeof field === 'string') { + return JSON.parse(field); + } + if (typeof field === 'object') { + return field; + } + throw new Error( + `expected a JSON string or an object, got ${typeof field}`); +} + class QueueEntry { /** @@ -27,14 +47,14 @@ class QueueEntry { } let entry; if (record.type === 'del') { - const overheadFields = record.overheadFields && JSON.parse(record.overheadFields); + const overheadFields = record.overheadFields && decodeField(record.overheadFields); entry = new DeleteOpQueueEntry(record.bucket, record.key, overheadFields); } else if (record.bucket === usersBucket) { // BucketQueueEntry class just handles puts of keys // to usersBucket entry = new BucketQueueEntry(record.key, record.value); } else if (record.value) { - const metadataVal = JSON.parse(record.value); + const metadataVal = decodeField(record.value); if (metadataVal.mdBucketModelVersion) { // it's bucket metadata entry = new BucketMdQueueEntry(record.key, metadataVal); diff --git a/tests/unit/lib/models/ObjectQueueEntry.spec.js b/tests/unit/lib/models/ObjectQueueEntry.spec.js index d9cfb96dc..5a682524a 100644 --- a/tests/unit/lib/models/ObjectQueueEntry.spec.js +++ b/tests/unit/lib/models/ObjectQueueEntry.spec.js @@ -384,5 +384,22 @@ describe('ObjectQueueEntry', () => { assert.strictEqual(entry.getKey(), 'key'); assert.deepStrictEqual(entry.getOverheadField('foo'), 'bar'); }); + + it('should create a DeleteOpQueueEntry with object overhead fields', + () => { + const entry = QueueEntry.createFromKafkaEntry({ + value: JSON.stringify({ + type: 'del', + bucket: 'bucket', + key: 'key', + overheadFields: { foo: 'bar' }, + }), + }); + + assert(entry instanceof DeleteOpQueueEntry); + assert.strictEqual(entry.getBucket(), 'bucket'); + assert.strictEqual(entry.getKey(), 'key'); + assert.deepStrictEqual(entry.getOverheadField('foo'), 'bar'); + }); }); }); diff --git a/tests/unit/replication/QueueEntry.spec.js b/tests/unit/lib/models/QueueEntry.spec.js similarity index 63% rename from tests/unit/replication/QueueEntry.spec.js rename to tests/unit/lib/models/QueueEntry.spec.js index 601c658f6..f430bc973 100644 --- a/tests/unit/replication/QueueEntry.spec.js +++ b/tests/unit/lib/models/QueueEntry.spec.js @@ -2,9 +2,8 @@ const assert = require('assert'); -const QueueEntry = - require('../../../lib/models/QueueEntry'); -const { replicationEntry } = require('../../utils/kafkaEntries'); +const QueueEntry = require('../../../../lib/models/QueueEntry'); +const { replicationEntry } = require('../../../utils/kafkaEntries'); describe('QueueEntry helper class', () => { describe('built from Kafka queue entry', () => { @@ -72,4 +71,54 @@ describe('QueueEntry helper class', () => { assert.strictEqual(completed1.getReplicationStatus(), 'COMPLETED'); }); }); + + describe('nested field decoding', () => { + function withDecodedValue(kafkaEntry) { + const record = JSON.parse(kafkaEntry.value); + record.value = JSON.parse(record.value); + return { ...kafkaEntry, value: JSON.stringify(record) }; + } + + it('should accept an object as the entry value', () => { + const fromString = QueueEntry.createFromKafkaEntry(replicationEntry); + const fromObject = QueueEntry.createFromKafkaEntry( + withDecodedValue(replicationEntry)); + + assert.strictEqual(fromObject.error, undefined); + assert.deepStrictEqual(fromObject.getValue(), fromString.getValue()); + assert.strictEqual(fromObject.getBucket(), fromString.getBucket()); + assert.strictEqual(fromObject.getObjectKey(), + fromString.getObjectKey()); + }); + + it('should report a malformed value string', () => { + const entry = QueueEntry.createFromKafkaEntry({ + value: JSON.stringify({ + type: 'put', + bucket: 'bucket', + key: 'key', + value: 'not json', + }), + }); + + assert.strictEqual(entry.error.message, + 'malformed JSON in kafka entry'); + }); + + it('should reject a value that is neither a string nor an object', + () => { + const entry = QueueEntry.createFromKafkaEntry({ + value: JSON.stringify({ + type: 'put', + bucket: 'bucket', + key: 'key', + value: 42, + }), + }); + + assert.strictEqual(entry.error.message, + 'malformed JSON in kafka entry'); + assert.match(entry.error.description, /got number/); + }); + }); });