]>
Commit | Line | Data |
---|---|---|
1 | import { forever, queue } from 'async' | |
2 | import * as Sequelize from 'sequelize' | |
3 | ||
4 | import { database as db } from '../../initializers/database' | |
5 | import { | |
6 | JOBS_FETCHING_INTERVAL, | |
7 | JOBS_FETCH_LIMIT_PER_CYCLE, | |
8 | JOB_STATES | |
9 | } from '../../initializers' | |
10 | import { logger } from '../../helpers' | |
11 | import { JobInstance } from '../../models' | |
12 | import { JobHandler, jobHandlers } from './handlers' | |
13 | ||
14 | type JobQueueCallback = (err: Error) => void | |
15 | ||
16 | class JobScheduler { | |
17 | ||
18 | private static instance: JobScheduler | |
19 | ||
20 | private constructor () { } | |
21 | ||
22 | static get Instance () { | |
23 | return this.instance || (this.instance = new this()) | |
24 | } | |
25 | ||
26 | activate () { | |
27 | const limit = JOBS_FETCH_LIMIT_PER_CYCLE | |
28 | ||
29 | logger.info('Jobs scheduler activated.') | |
30 | ||
31 | const jobsQueue = queue<JobInstance, JobQueueCallback>(this.processJob.bind(this)) | |
32 | ||
33 | // Finish processing jobs from a previous start | |
34 | const state = JOB_STATES.PROCESSING | |
35 | db.Job.listWithLimit(limit, state) | |
36 | .then(jobs => { | |
37 | this.enqueueJobs(jobsQueue, jobs) | |
38 | ||
39 | forever( | |
40 | next => { | |
41 | if (jobsQueue.length() !== 0) { | |
42 | // Finish processing the queue first | |
43 | return setTimeout(next, JOBS_FETCHING_INTERVAL) | |
44 | } | |
45 | ||
46 | const state = JOB_STATES.PENDING | |
47 | db.Job.listWithLimit(limit, state) | |
48 | .then(jobs => { | |
49 | this.enqueueJobs(jobsQueue, jobs) | |
50 | ||
51 | // Optimization: we could use "drain" from queue object | |
52 | return setTimeout(next, JOBS_FETCHING_INTERVAL) | |
53 | }) | |
54 | .catch(err => logger.error('Cannot list pending jobs.', err)) | |
55 | }, | |
56 | ||
57 | err => logger.error('Error in job scheduler queue.', err) | |
58 | ) | |
59 | }) | |
60 | .catch(err => logger.error('Cannot list pending jobs.', err)) | |
61 | } | |
62 | ||
63 | createJob (transaction: Sequelize.Transaction, handlerName: string, handlerInputData: object) { | |
64 | const createQuery = { | |
65 | state: JOB_STATES.PENDING, | |
66 | handlerName, | |
67 | handlerInputData | |
68 | } | |
69 | const options = { transaction } | |
70 | ||
71 | return db.Job.create(createQuery, options) | |
72 | } | |
73 | ||
74 | private enqueueJobs (jobsQueue: AsyncQueue<JobInstance>, jobs: JobInstance[]) { | |
75 | jobs.forEach(job => jobsQueue.push(job)) | |
76 | } | |
77 | ||
78 | private processJob (job: JobInstance, callback: (err: Error) => void) { | |
79 | const jobHandler = jobHandlers[job.handlerName] | |
80 | if (jobHandler === undefined) { | |
81 | logger.error('Unknown job handler for job %s.', job.handlerName) | |
82 | return callback(null) | |
83 | } | |
84 | ||
85 | logger.info('Processing job %d with handler %s.', job.id, job.handlerName) | |
86 | ||
87 | job.state = JOB_STATES.PROCESSING | |
88 | return job.save() | |
89 | .then(() => { | |
90 | return jobHandler.process(job.handlerInputData) | |
91 | }) | |
92 | .then( | |
93 | result => { | |
94 | return this.onJobSuccess(jobHandler, job, result) | |
95 | }, | |
96 | ||
97 | err => { | |
98 | logger.error('Error in job handler %s.', job.handlerName, err) | |
99 | return this.onJobError(jobHandler, job, err) | |
100 | } | |
101 | ) | |
102 | .then(() => callback(null)) | |
103 | .catch(err => { | |
104 | this.cannotSaveJobError(err) | |
105 | return callback(err) | |
106 | }) | |
107 | } | |
108 | ||
109 | private onJobError (jobHandler: JobHandler<any>, job: JobInstance, err: Error) { | |
110 | job.state = JOB_STATES.ERROR | |
111 | ||
112 | return job.save() | |
113 | .then(() => jobHandler.onError(err, job.id)) | |
114 | .catch(err => this.cannotSaveJobError(err)) | |
115 | } | |
116 | ||
117 | private onJobSuccess (jobHandler: JobHandler<any>, job: JobInstance, jobResult: any) { | |
118 | job.state = JOB_STATES.SUCCESS | |
119 | ||
120 | return job.save() | |
121 | .then(() => jobHandler.onSuccess(job.id, jobResult)) | |
122 | .catch(err => this.cannotSaveJobError(err)) | |
123 | } | |
124 | ||
125 | private cannotSaveJobError (err: Error) { | |
126 | logger.error('Cannot save new job state.', err) | |
127 | } | |
128 | } | |
129 | ||
130 | // --------------------------------------------------------------------------- | |
131 | ||
132 | export { | |
133 | JobScheduler | |
134 | } |