2015-12-22 02:57:10 +08:00
|
|
|
'use strict';
|
|
|
|
|
2016-01-23 01:22:21 +08:00
|
|
|
var errorCodes = require('../app/postgresql/error_codes').codeToCondition;
|
2015-12-22 02:57:10 +08:00
|
|
|
var PSQL = require('cartodb-psql');
|
2016-03-18 21:57:18 +08:00
|
|
|
var queue = require('queue-async');
|
2015-12-22 02:57:10 +08:00
|
|
|
|
2016-01-22 19:43:41 +08:00
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
function JobRunner(jobBackend, userDatabaseMetadataService) {
|
|
|
|
this.jobBackend = jobBackend;
|
2015-12-22 02:57:10 +08:00
|
|
|
this.userDatabaseMetadataService = userDatabaseMetadataService;
|
|
|
|
}
|
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
JobRunner.prototype.run = function (job_id, callback) {
|
2015-12-22 02:57:10 +08:00
|
|
|
var self = this;
|
2016-01-08 18:32:01 +08:00
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
self.jobBackend.get(job_id, function (err, job) {
|
2015-12-22 02:57:10 +08:00
|
|
|
if (err) {
|
2016-01-08 22:47:59 +08:00
|
|
|
return callback(err);
|
2015-12-22 02:57:10 +08:00
|
|
|
}
|
2015-12-29 17:19:10 +08:00
|
|
|
|
2015-12-31 03:16:18 +08:00
|
|
|
if (job.status !== 'pending') {
|
2016-01-13 23:25:25 +08:00
|
|
|
var error = new Error('Cannot run job ' + job.job_id + ' due to its status is ' + job.status);
|
|
|
|
error.name = 'InvalidJobStatus';
|
|
|
|
return callback(error);
|
2015-12-31 03:16:18 +08:00
|
|
|
}
|
|
|
|
|
2015-12-22 02:57:10 +08:00
|
|
|
self.userDatabaseMetadataService.getUserMetadata(job.user, function (err, userDatabaseMetadata) {
|
|
|
|
if (err) {
|
2016-01-08 22:47:59 +08:00
|
|
|
return callback(err);
|
2015-12-22 02:57:10 +08:00
|
|
|
}
|
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
self.jobBackend.setRunning(job, function (err, job) {
|
|
|
|
if (err) {
|
|
|
|
return callback(err);
|
|
|
|
}
|
2015-12-22 02:57:10 +08:00
|
|
|
|
2016-03-18 21:57:18 +08:00
|
|
|
self._series(job, userDatabaseMetadata, callback);
|
2016-01-08 22:47:59 +08:00
|
|
|
});
|
|
|
|
});
|
|
|
|
});
|
|
|
|
};
|
2015-12-22 02:57:10 +08:00
|
|
|
|
2016-03-18 21:57:18 +08:00
|
|
|
JobRunner.prototype._series = function(job, userDatabaseMetadata, callback) {
|
|
|
|
var jobQueue = queue(1); // performs in series
|
|
|
|
|
|
|
|
if (!Array.isArray(job.query)) {
|
|
|
|
job.query = [ job.query ];
|
|
|
|
}
|
|
|
|
|
|
|
|
for (var i = 0; i < job.query.length; i++) {
|
|
|
|
jobQueue.defer(this._query.bind(this), job, userDatabaseMetadata, i);
|
|
|
|
}
|
|
|
|
|
|
|
|
jobQueue.await(function (err, result) {
|
|
|
|
if (err) {
|
|
|
|
return callback(err);
|
|
|
|
}
|
|
|
|
|
|
|
|
// last result is the good one
|
|
|
|
if (Array.isArray(result)) {
|
|
|
|
return callback(null, result[result.length - 1]);
|
|
|
|
}
|
|
|
|
|
|
|
|
callback(null, result);
|
|
|
|
})
|
|
|
|
};
|
|
|
|
|
|
|
|
JobRunner.prototype._query = function (job, userDatabaseMetadata, index, callback) {
|
2016-01-08 22:47:59 +08:00
|
|
|
var self = this;
|
|
|
|
|
|
|
|
var pg = new PSQL(userDatabaseMetadata, {}, { destroyOnError: true });
|
2015-12-22 22:43:00 +08:00
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
pg.query('SET statement_timeout=0', function (err) {
|
|
|
|
if(err) {
|
|
|
|
return self.jobBackend.setFailed(job, err, callback);
|
|
|
|
}
|
2015-12-31 03:16:18 +08:00
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
// mark query to allow to users cancel their queries whether users request for it
|
2016-03-18 21:57:18 +08:00
|
|
|
var sql = job.query[index] + ' /* ' + job.job_id + ' */';
|
2015-12-24 00:29:11 +08:00
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
pg.eventedQuery(sql, function (err, query) {
|
|
|
|
if (err) {
|
|
|
|
return self.jobBackend.setFailed(job, err, callback);
|
|
|
|
}
|
2015-12-31 22:42:31 +08:00
|
|
|
|
2016-01-08 22:47:59 +08:00
|
|
|
query.on('error', function (err) {
|
2016-01-18 02:28:15 +08:00
|
|
|
// if query has been cancelled then it's going to get the current job status saved by query_canceller
|
2016-01-23 01:22:21 +08:00
|
|
|
if (errorCodes[err.code.toString()] === 'query_canceled') {
|
2016-01-18 02:28:15 +08:00
|
|
|
return self.jobBackend.get(job.job_id, callback);
|
2016-01-08 22:47:59 +08:00
|
|
|
}
|
2016-01-18 02:28:15 +08:00
|
|
|
|
|
|
|
self.jobBackend.setFailed(job, err, callback);
|
2016-01-08 22:47:59 +08:00
|
|
|
});
|
|
|
|
|
|
|
|
query.on('end', function (result) {
|
2016-01-18 02:28:15 +08:00
|
|
|
// only if result is present then query is done sucessfully otherwise an error has happened
|
|
|
|
// and it was handled by error listener
|
2016-01-08 22:47:59 +08:00
|
|
|
if (result) {
|
2016-01-18 02:28:15 +08:00
|
|
|
return self.jobBackend.setDone(job, callback);
|
2016-01-08 22:47:59 +08:00
|
|
|
}
|
2015-12-22 02:57:10 +08:00
|
|
|
});
|
|
|
|
});
|
|
|
|
});
|
|
|
|
};
|
|
|
|
|
|
|
|
module.exports = JobRunner;
|