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
18 changes: 13 additions & 5 deletions extensions/replication/queueProcessor/QueueProcessor.js
Original file line number Diff line number Diff line change
Expand Up @@ -264,11 +264,18 @@ class QueueProcessor extends EventEmitter {
}
}

this.taskScheduler = new TaskScheduler(
// The two consumers' entries must not deduplicate against each other
this.taskScheduler = this._createTaskScheduler();
this.dataMoverTaskScheduler = this._createTaskScheduler();
}

_createTaskScheduler() {
return new TaskScheduler(
(ctx, done) => ctx.task.processQueueEntry(
ctx.entry, ctx.kafkaEntry, done),
ctx => getTaskSchedulerQueueKey(ctx.entry),
ctx => getTaskSchedulerDedupeKey(ctx.entry));
ctx => getTaskSchedulerDedupeKey(ctx.entry),
this.logger);
}

_setupVaultclientCache() {
Expand Down Expand Up @@ -989,9 +996,10 @@ class QueueProcessor extends EventEmitter {
this.logger.debug('data mover entry is being pushed', {
entry: actionEntry.getLogInfo(),
});
return this.taskScheduler.push({ task, entry: actionEntry,
kafkaEntry },
done);
return this.dataMoverTaskScheduler.push(
{ task, entry: actionEntry, kafkaEntry },
done,
);
}
this.logger.debug('skip data mover entry', {
entry: actionEntry.getLogInfo(),
Expand Down
10 changes: 5 additions & 5 deletions extensions/replication/queueProcessor/taskSchedulerHelpers.js
Comment thread
delthas marked this conversation as resolved.
Comment thread
delthas marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -14,15 +14,15 @@ function getTaskSchedulerQueueKey(entry) {

function getTaskSchedulerDedupeKey(entry) {
Comment thread
delthas marked this conversation as resolved.
if (entry instanceof ObjectQueueEntry) {
const key = entry.getCanonicalKey();
const canonicalKey = entry.getCanonicalKey();
const version = entry.getVersionId();
const contentMd5 = entry.getContentMd5();
return `${key}:${version || ''}:${contentMd5}`;
return `${canonicalKey}:${version || ''}:${contentMd5}`;
}
if (entry instanceof ActionQueueEntry) {
const { key, version, contentMd5 } =
entry.getAttribute('target');
return `${key}:${version || ''}:${contentMd5}`;
const { bucket, key, version, eTag } = entry.getAttribute('target');
const objectETag = (eTag || '').replace(/"/g, '');
return `${bucket}/${key}:${version || ''}:${objectETag}`;
}
return undefined;
}
Expand Down
2 changes: 1 addition & 1 deletion lib/BackbeatConsumer.js
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ class BackbeatConsumer extends EventEmitter {
(ctx, done) => this._processTask(ctx.entry, done),
orderByFunc,
null,
this._concurrency,
Comment thread
delthas marked this conversation as resolved.
this._log,
);

// Single drain callback fed by both signals (processing queue
Expand Down
8 changes: 7 additions & 1 deletion lib/tasks/TaskScheduler.js
Original file line number Diff line number Diff line change
Expand Up @@ -29,11 +29,14 @@ class TaskScheduler {
* @param {function} [getDedupeKeyFunc] - function to get the dedupe
* key of a task: getDedupeKeyFunc(ctx) -> {string} key
* (see: {@link TaskScheduler.push()})
* @param {Logger} [logger] - logger object, used to trace the tasks
* skipped by deduplication
*/
constructor(processingFunc, getQueueKeyFunc, getDedupeKeyFunc) {
constructor(processingFunc, getQueueKeyFunc, getDedupeKeyFunc, logger) {
this._processingFunc = processingFunc;
this._getQueueKeyFunc = getQueueKeyFunc;
this._getDedupeKeyFunc = getDedupeKeyFunc;
this._logger = logger;
this._drainFunc = null;
this._taskQueues = {};
this._dedupeCache = {};
Expand Down Expand Up @@ -87,6 +90,9 @@ class TaskScheduler {
}
if (dedupeKey !== undefined) {
if (this._dedupeCache[dedupeKey]) {
this._logger?.debug('task skipped: a task with the ' +
'same dedupe key is in progress ' +
'or queued', { dedupeKey });
return process.nextTick(done);
}
this._dedupeCache[dedupeKey] = true;
Expand Down
24 changes: 24 additions & 0 deletions tests/unit/lib/tasks/TaskScheduler.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -360,6 +360,30 @@ describe('TaskScheduler', () => {
}, 10);
});
});
it('should log the tasks skipped by deduplication', done => {
let blockedTaskResolve;
const blockedTaskPromise = new Promise(resolve => { blockedTaskResolve = resolve; });
const processingFunc = sinon.stub().callsFake((ctx, cb) => {
if (ctx.id === 'block') {
return blockedTaskPromise.then(() => cb());
}
return cb();
});
const logger = { debug: sinon.spy() };
const taskScheduler = new TaskScheduler(
processingFunc, null, ctx => ctx.dedupeKey, logger);

taskScheduler.push({ id: 'block', dedupeKey: 'key' }, () => {});
taskScheduler.push({ dedupeKey: 'key' }, () => {});

process.nextTick(() => {
assert(logger.debug.calledOnce);
assert.deepStrictEqual(logger.debug.firstCall.args[1],
{ dedupeKey: 'key' });
blockedTaskResolve();
done();
});
});
it('should remove dedup key after task completion', done => {
const processingFunc = sinon.stub().yields();
const taskScheduler = new TaskScheduler(processingFunc, ctx => ctx.queueKey, null);
Expand Down
30 changes: 30 additions & 0 deletions tests/unit/replication/QueueProcessor.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -305,6 +305,36 @@ describe('Queue Processor', () => {
});
});

describe('processDataMoverEntry', () => {
beforeEach(() => {
qp.taskScheduler = { push: sinon.stub().callsArgWith(1, null) };
qp.dataMoverTaskScheduler = {
push: sinon.stub().callsArgWith(1, null),
};
qp._mProducer = { getProducer: () => null };
});

it('dispatches copyLocation actions to the data mover scheduler',
done => {
const kafkaEntry = {
value: JSON.stringify({
action: 'copyLocation',
toLocation: 'site-crr',
target: {
bucket: 'src-bucket',
key: 'obj',
eTag: '"d41d8cd98f00b204e9800998ecf8427e"',
},
}),
};
qp.processDataMoverEntry(kafkaEntry, () => {
sinon.assert.calledOnce(qp.dataMoverTaskScheduler.push);
sinon.assert.notCalled(qp.taskScheduler.push);
done();
});
});
});

describe('constructor', () => {
it('should use s3c site\'s host as a destination host', () => {
const config = getQueueProcessorConfig();
Expand Down
111 changes: 101 additions & 10 deletions tests/unit/replication/taskSchedulerHelpers.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,54 +7,63 @@ const { ObjectMD } = require('arsenal').models;
const { getTaskSchedulerQueueKey, getTaskSchedulerDedupeKey } = require(
'../../../extensions/replication/queueProcessor/taskSchedulerHelpers');

function makeObjectQueueEntry({ key, versionId, contentMd5 }) {
function makeObjectQueueEntry({ bucket, key, versionId, contentMd5 }) {
const objMd = new ObjectMD().setKey(key);
return new ObjectQueueEntry('test-bucket-name', key, objMd)
return new ObjectQueueEntry(bucket, key, objMd)
.setVersionId(versionId)
.setContentMd5(contentMd5);
}

function makeActionQueueEntry({ key, versionId, contentMd5 }) {
function makeActionQueueEntry({ bucket, key, versionId, eTag }) {
return new ActionQueueEntry({
action: 'copyLocation',
target: {
bucket,
key,
version: versionId,
contentMd5,
eTag,
}
});
}

function makeQueueEntry(entryClass, { key, versionId, contentMd5 }) {
function makeQueueEntry(entryClass, { bucket, key, versionId, contentMd5 }) {
if (entryClass === 'ObjectQueueEntry') {
return makeObjectQueueEntry({ key, versionId, contentMd5 });
return makeObjectQueueEntry({ bucket, key, versionId, contentMd5 });
}
if (entryClass === 'ActionQueueEntry') {
return makeActionQueueEntry({ key, versionId, contentMd5 });
return makeActionQueueEntry({ bucket, key, versionId,
eTag: `"${contentMd5}"` });
}
return assert.fail(`bad class ${entryClass}`);
}

function makeTestEntry(entryClass,
{ keySelector, versionIdSelector, contentMd5Selector }) {
{ bucketSelector, keySelector, versionIdSelector,
contentMd5Selector }) {
const buckets = ['test-bucket-name', 'other-bucket-name'];
const keys = ['masterkey1', 'masterkey2'];
const versionIds = ['abcdef', 'ghijkl'];
const contentMd5s = ['d41d8cd98f00b204e9800998ecf8427e',
'93b07384d113edec49eaa6238ad5ff00'];
return makeQueueEntry(
entryClass, {
bucket: buckets[bucketSelector],
key: keys[keySelector],
versionId: versionIds[versionIdSelector],
contentMd5: contentMd5s[contentMd5Selector],
});
}

function makeTestEntryPair(entryClass, { distinctKey, distinctVersionId,
function makeTestEntryPair(entryClass, { distinctBucket, distinctKey,
distinctVersionId,
distinctContentMd5 }) {
return [
makeTestEntry(entryClass, { keySelector: 0,
makeTestEntry(entryClass, { bucketSelector: 0,
keySelector: 0,
versionIdSelector: 0,
contentMd5Selector: 0 }),
makeTestEntry(entryClass, {
bucketSelector: distinctBucket ? 1 : 0,
keySelector: distinctKey ? 1 : 0,
versionIdSelector: distinctVersionId ? 1 : 0,
contentMd5Selector: distinctContentMd5 ? 1 : 0,
Expand Down Expand Up @@ -95,6 +104,15 @@ describe('QueueProcessor::getTaskSchedulerQueueKey', () => {
});
assert.notStrictEqual(queueKey1, queueKey2);
});

it(`should return different keys of ${entryClass} with the same ` +
'master key in different buckets', () => {
const [queueKey1, queueKey2] = makeQueueKeyPair(
entryClass, {
distinctBucket: true,
});
assert.notStrictEqual(queueKey1, queueKey2);
});
});
});

Expand Down Expand Up @@ -134,5 +152,78 @@ describe('QueueProcessor::getTaskSchedulerDedupeKey', () => {
});
assert.notStrictEqual(dedupeKey1, dedupeKey2);
});

it(`should return different keys of ${entryClass} with the same ` +
'master-key/version/md5 in different buckets', () => {
const [dedupeKey1, dedupeKey2] = makeDedupeKeyPair(
entryClass, {
distinctBucket: true,
});
assert.notStrictEqual(dedupeKey1, dedupeKey2);
});
});

});

describe('QueueProcessor::getTaskSchedulerDedupeKey of copyLocation actions',
() => {
const objectParams = {
bucket: 'test-bucket-name',
key: 'index.html',
eTag: '"d41d8cd98f00b204e9800998ecf8427e"',
};

it('should return different keys for the same key of non-versioned ' +
'objects in different buckets', () => {
const entry1 = makeActionQueueEntry(objectParams);
const entry2 = makeActionQueueEntry(
Object.assign({}, objectParams, { bucket: 'other-bucket-name' }));
assert.notStrictEqual(getTaskSchedulerDedupeKey(entry1),
getTaskSchedulerDedupeKey(entry2));
});

it('should return different keys for a non-versioned object ' +
'overwritten with new contents', () => {
const entry1 = makeActionQueueEntry(objectParams);
const entry2 = makeActionQueueEntry(
Object.assign({}, objectParams,
{ eTag: '"93b07384d113edec49eaa6238ad5ff00"' }));
assert.notStrictEqual(getTaskSchedulerDedupeKey(entry1),
getTaskSchedulerDedupeKey(entry2));
});

it('should return matching keys for duplicates of the same action', () => {
assert.strictEqual(
getTaskSchedulerDedupeKey(makeActionQueueEntry(objectParams)),
getTaskSchedulerDedupeKey(makeActionQueueEntry(objectParams)));
});

it('should return matching keys whether the eTag is quoted or not',
() => {
const quoted = makeActionQueueEntry(objectParams);
const unquoted = makeActionQueueEntry(
Object.assign({}, objectParams,
{ eTag: 'd41d8cd98f00b204e9800998ecf8427e' }));
assert.strictEqual(getTaskSchedulerDedupeKey(quoted),
getTaskSchedulerDedupeKey(unquoted));
});

it('should return matching keys for actions without an eTag', () => {
const params = Object.assign({}, objectParams, { eTag: undefined });
assert.strictEqual(
getTaskSchedulerDedupeKey(makeActionQueueEntry(params)),
getTaskSchedulerDedupeKey(makeActionQueueEntry(params)));
});

it('should return different keys for MPU objects differing only by ' +
'their part count', () => {
const entry1 = makeActionQueueEntry(
Object.assign({}, objectParams,
{ eTag: '"d41d8cd98f00b204e9800998ecf8427e-2"' }));
const entry2 = makeActionQueueEntry(
Object.assign({}, objectParams,
{ eTag: '"d41d8cd98f00b204e9800998ecf8427e-3"' }));
assert.notStrictEqual(getTaskSchedulerDedupeKey(entry1),
getTaskSchedulerDedupeKey(entry2));
});
});
Loading