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