2015-12-29 17:19:10 +08:00
|
|
|
'use strict';
|
|
|
|
|
2016-10-18 00:16:52 +08:00
|
|
|
var debug = require('./util/debug')('queue');
|
|
|
|
|
2017-03-31 20:30:33 +08:00
|
|
|
function JobQueue(metadataBackend, jobPublisher, queueIndex) {
|
2015-12-29 17:19:10 +08:00
|
|
|
this.metadataBackend = metadataBackend;
|
2016-06-30 00:29:53 +08:00
|
|
|
this.jobPublisher = jobPublisher;
|
2017-03-31 20:30:33 +08:00
|
|
|
this.queueIndex = queueIndex;
|
2015-12-29 17:19:10 +08:00
|
|
|
}
|
|
|
|
|
2016-10-12 23:53:03 +08:00
|
|
|
module.exports = JobQueue;
|
|
|
|
|
|
|
|
var QUEUE = {
|
|
|
|
DB: 5,
|
2017-03-31 20:30:33 +08:00
|
|
|
PREFIX: 'batch:queue:',
|
|
|
|
INDEX: 'batch:indexes:queue'
|
2016-10-12 23:53:03 +08:00
|
|
|
};
|
2017-03-31 20:30:33 +08:00
|
|
|
|
2016-10-12 23:53:03 +08:00
|
|
|
module.exports.QUEUE = QUEUE;
|
|
|
|
|
2016-10-13 03:32:29 +08:00
|
|
|
JobQueue.prototype.enqueue = function (user, jobId, callback) {
|
2016-10-18 00:16:52 +08:00
|
|
|
debug('JobQueue.enqueue user=%s, jobId=%s', user, jobId);
|
2017-03-31 20:30:33 +08:00
|
|
|
|
|
|
|
this.metadataBackend.redisMultiCmd(QUEUE.DB, [
|
|
|
|
[ 'LPUSH', QUEUE.PREFIX + user, jobId ],
|
|
|
|
[ 'SADD', QUEUE.INDEX, user ]
|
|
|
|
], function (err) {
|
2016-06-30 00:29:53 +08:00
|
|
|
if (err) {
|
|
|
|
return callback(err);
|
|
|
|
}
|
|
|
|
|
2016-10-18 00:16:52 +08:00
|
|
|
this.jobPublisher.publish(user);
|
2016-06-30 00:29:53 +08:00
|
|
|
callback();
|
2016-10-18 00:16:52 +08:00
|
|
|
}.bind(this));
|
2015-12-29 17:19:10 +08:00
|
|
|
};
|
|
|
|
|
2016-10-13 04:40:09 +08:00
|
|
|
JobQueue.prototype.size = function (user, callback) {
|
|
|
|
this.metadataBackend.redisCmd(QUEUE.DB, 'LLEN', [ QUEUE.PREFIX + user ], callback);
|
|
|
|
};
|
|
|
|
|
2016-10-13 03:32:29 +08:00
|
|
|
JobQueue.prototype.dequeue = function (user, callback) {
|
2017-03-31 20:30:33 +08:00
|
|
|
var dequeueScript = [
|
|
|
|
'local job_id = redis.call("RPOP", KEYS[1])',
|
|
|
|
'if redis.call("LLEN", KEYS[1]) == 0 then',
|
|
|
|
' redis.call("SREM", KEYS[2], ARGV[1])',
|
|
|
|
'end',
|
|
|
|
'return job_id'
|
|
|
|
].join('\n');
|
|
|
|
|
|
|
|
var redisParams = [
|
|
|
|
dequeueScript, //lua source code
|
|
|
|
2, // Two "keys" to pass
|
|
|
|
QUEUE.PREFIX + user, //KEYS[1], the key of the queue
|
|
|
|
QUEUE.INDEX, //KEYS[2], the key of the index
|
|
|
|
user // ARGV[1] - value of the element to remove form the index
|
|
|
|
];
|
|
|
|
|
|
|
|
this.metadataBackend.redisCmd(QUEUE.DB, 'EVAL', redisParams, function (err, jobId) {
|
2016-10-18 00:16:52 +08:00
|
|
|
debug('JobQueue.dequeued user=%s, jobId=%s', user, jobId);
|
|
|
|
return callback(err, jobId);
|
|
|
|
});
|
2015-12-29 17:19:10 +08:00
|
|
|
};
|
|
|
|
|
2016-10-13 03:32:29 +08:00
|
|
|
JobQueue.prototype.enqueueFirst = function (user, jobId, callback) {
|
2016-10-18 00:16:52 +08:00
|
|
|
debug('JobQueue.enqueueFirst user=%s, jobId=%s', user, jobId);
|
2016-10-13 03:32:29 +08:00
|
|
|
this.metadataBackend.redisCmd(QUEUE.DB, 'RPUSH', [ QUEUE.PREFIX + user, jobId ], callback);
|
2016-01-13 23:25:25 +08:00
|
|
|
};
|