Repository navigation
Expand file tree
/
Copy pathqueue.spec.ts
More file actions
117 lines (95 loc) · 2.37 KB
/
Copy pathqueue.spec.ts
File metadata and controls
117 lines (95 loc) · 2.37 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
import { Queue, Job } from './queue'
import { expect } from 'chai'
import knex from 'knex'
import 'mocha'
import { setupTables } from './migrations'
const db = knex({
client: 'postgresql',
connection: {
database: 'postqueue_test' || process.env.POSTGRES_DB,
user: 'postgres' || process.env.POSTRGES_USER,
password: '' || process.env.POSTGRES_PASS
}
})
describe('queue', () => {
before(async () => {
const tablesExist = await db.schema.hasTable('postqueue')
if (!tablesExist) {
await setupTables(db)
}
})
it('should successfully add and delete', async () => {
const q = new Queue('test-queue', db)
const job = await q.add({
test: 'hello-world'
})
await job.remove()
})
it('should process jobs', (done) => {
(async () => {
const q = new Queue('test-queue', db)
await q.add({
test: 'hello-world'
})
q.process(async (j) => {
expect(j.data).to.not.be.equal(null)
expect(j.data.test).to.be.equal('hello-world')
q.shutdown()
await j.remove()
done()
return null
})
})()
})
it('should process recurring jobs multiple times', (done) => {
(async () => {
const q = new Queue('test-queue', db)
let triggered = 0
const job = await q.add({
test: 'hello-world'
}, {
everySecs: 1
})
q.process(async (j: Job) => {
expect(j.data).to.not.be.equal(null)
expect(j.data.test).to.be.equal('hello-world')
triggered++
if (triggered === 2) {
q.shutdown()
await j.remove()
done()
}
return null
})
})()
}).timeout(5000)
it('should continue processing after error', (done) => {
(async () => {
const q = new Queue('test-queue', db)
let triggered = 0
for (let i = 0; i < 100; i++) {
await q.add({
test: 'hello-world'
})
}
q.process(async (j: Job) => {
expect(j.data).to.not.be.equal(null)
expect(j.data.test).to.be.equal('hello-world')
triggered++
if (triggered === 20) {
throw new Error('unhandled error')
}
if (triggered === 100) {
q.shutdown()
done()
}
return null
})
})()
})
after(() => {
setImmediate(() => {
db.destroy()
})
})
})