forked from OptimalBits/bull
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmessage.js
More file actions
125 lines (110 loc) · 2.76 KB
/
Copy pathmessage.js
File metadata and controls
125 lines (110 loc) · 2.76 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
"use strict";
var redis = require('redis');
var when = require('when');
/**
interface JobOptions
{
attempts: number;
}
*/
// queue: Queue, msgId: string, data: {}, opts: JobOptions
var Message = function Message(queue, msgId, data, opts){
this.queue = queue;
this.msgId = msgId;
this.data = data;
this.opts = opts;
this._progress = 0;
}
Message.create = function(queue, msgId, data, opts){
var deferred = when.defer();
var msg = new Message(queue, msgId, data, opts);
queue.client.HMSET(queue.toKey(msgId), msg.toData(), function(err){
if(err){
deferred.reject(err);
}else{
deferred.resolve(job);
}
});
return deferred.promise;
}
Message.fromId = function(queue, msgId){
var deferred = when.defer();
queue.client.HGETALL(queue.toKey(msgId), function(err, data){
if(data){
deferred.resolve(Message.fromData(queue, msgId, data));
}else{
deferred.reject(err);
}
});
return deferred.promise;
}
Message.prototype.toData = function(){
return {
name: this.name,
data: JSON.stringify(this.data || {}),
opts: JSON.stringify(this.opts || {}),
progress: this._progress
}
}
Message.prototype.progress = function(progress){
if(progress){
var deferred = when.defer();
var _this = this;
this.queue.client.hset(this.queue.toKey(this.msgId), 'progress', progress, function(err){
if(err){
deferred.reject(err);
}else{
deferred.resolve();
_this.queue.emit('progress', _this, progress);
}
});
return deferred.promise;
}else{
return this._progress;
}
}
Job.prototype.completed = function(){
return this._done('completed');
}
Job.prototype.failed = function(err){
return this._done('failed');
}
Job.prototype.isCompleted = function(){
return this._isDone('completed');
}
Job.prototype.isFailed = function(){
return this._isDone('failed');
}
Job.prototype._isDone = function(list){
var deferred = when.defer();
this.queue.client.SISMEMBER(this.queue.toKey(list), this.jobId, function(err, isMember){
if(err){
deferred.reject(err);
}else{
deferred.resolve(isMember === 1);
}
});
return deferred.promise;
}
Job.prototype._done = function(list){
var deferred = when.defer();
var queue = this.queue;
var activeList = queue.toKey('active');
var completedList = queue.toKey(list);
queue.client.multi()
.lrem(activeList, 0, this.jobId)
.sadd(completedList, this.jobId)
.exec(function(err){
!err && deferred.resolve();
err && deferred.reject(err);
});
return deferred.promise;
}
/**
*/
Job.fromData = function(queue, jobId, data){
var job = new Job(queue, jobId, data.name, JSON.parse(data.data), data.opts);
job._progress = parseInt(data.progress);
return job;
}
module.exports = Job;