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
24 changes: 22 additions & 2 deletions lib/models/QueueEntry.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

/**
Expand All @@ -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);
Expand Down
17 changes: 17 additions & 0 deletions tests/unit/lib/models/ObjectQueueEntry.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -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');
});
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -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', () => {
Expand Down Expand Up @@ -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/);
});
});
});
Loading