forked from Seneca-CDOT/telescope
-
Notifications
You must be signed in to change notification settings - Fork 0
/
queue.js
97 lines (92 loc) · 3.01 KB
/
queue.js
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
const Bull = require('bull');
const { createRedisClient } = require('./redis');
const { logger } = require('../utils/logger');
/**
* Shared redis connections for pub/sub, see:
* https://github.com/OptimalBits/bull/blob/28a2b9aa444d028fc5192c9bbdc9bb5811e77b08/PATTERNS.md#reusing-redis-connections
*/
const client = createRedisClient();
const subscriber = createRedisClient();
/**
* Create a Queue with the given `name` (String).
* We create a Bull Queue using either a real or mocked
* redis, and manage the creation of the redis connections.
* We also setup logging for this queue name.
*/
function createQueue(name) {
const queue = new Bull(name, {
createClient: type => {
switch (type) {
case 'client':
return client;
case 'subscriber':
return subscriber;
default:
return createRedisClient();
}
},
})
.on('error', err => {
// An error occurred
if (err.code === 'ECONNREFUSED') {
logger.info(
'\n\n\t💡 It appears that Redis is not running on your machine.',
'\n\t Please see our documentation for how to install and run Redis:',
'\n\t https://github.com/Seneca-CDOT/telescope/blob/master/docs/CONTRIBUTING.md\n'
);
} else {
logger.error({ err }, `Queue ${name} error`);
}
})
.on('waiting', jobID => {
// A job is waiting for the next idling worker
logger.debug(`Job ${jobID} is waiting.`);
})
.on('active', job => {
// A job has started (use jobPromise.cancel() to abort it)
logger.debug(`Job ${job.id} is active`);
})
.on('stalled', job => {
// A job was marked as stalled. This is useful for debugging
// which workers are crashing or pausing the event loop
logger.debug(`Job ${job.id} has stalled.`);
})
.on('progress', (job, progress) => {
// A job's progress was updated
logger.debug(`Job ${job.id} progress:`, progress);
})
.on('completed', job => {
// A job has been completed
logger.debug(`Job ${job.id} completed.`);
})
.on('failed', (job, error) => {
// A job failed with an error
logger.error({ error }, `Job ${job.id} failed.`);
})
.on('paused', job => {
// The queue was paused
logger.debug(`Queue ${name} resumed. ID:`, job.id);
})
.on('resumed', job => {
// The queue resumed
logger.debug(`Queue ${name} resumed. ID: `, job.id);
})
.on('cleaned', (jobs, types) => {
// Old jobs were cleaned from the queue
// 'Jobs' is an array of cleaned jobs
// 'Types' is an array of their types
logger.debug(`Queue ${name} was cleaned. Jobs: `, jobs, ' Types: ', types);
})
.on('drained', () => {
// The queue was drained
// (the last item in the queue was returned by a worker)
logger.debug(`Queue ${name} was drained.`);
})
.on('removed', job => {
logger.debug(`Job ${job.id} was removed.`);
});
return queue;
}
module.exports = {
createQueue,
};