nci/distributor.js

132 lines
3.2 KiB
JavaScript
Raw Normal View History

'use strict';
var Steppy = require('twostep').Steppy,
_ = require('underscore'),
Distributor = require('./lib/distributor').Distributor,
getAvgProjectBuildDuration = (
require('./lib/project').getAvgProjectBuildDuration
),
db = require('./db'),
2015-10-04 16:57:53 +00:00
logger = require('./lib/logger')('distributor');
exports.init = function(app, callback) {
var distributor = new Distributor({
nodes: app.config.nodes,
projects: app.projects,
saveBuild: function(build, callback) {
Steppy(
function() {
2015-11-19 17:28:51 +00:00
if (_(build.project).has('avgBuildDuration')) {
this.pass(build.project.avgBuildDuration);
2015-11-19 17:28:51 +00:00
} else {
getAvgProjectBuildDuration(build.project.name, this.slot());
}
},
function(err, avgBuildDuration) {
build.project.avgBuildDuration = avgBuildDuration;
db.builds.put(build, this.slot());
},
function() {
this.pass(build);
},
callback
);
}
});
2015-11-19 17:28:51 +00:00
var buildDataResourcesHash = {};
// create resource for build data
var createBuildDataResource = function(buildId) {
if (buildId in buildDataResourcesHash) {
return;
}
var buildDataResource = app.dataio.resource('build' + buildId);
buildDataResource.on('connection', function(client) {
2015-11-18 20:14:46 +00:00
var callback = this.async();
Steppy(
function() {
db.logLines.find({
2015-11-19 16:07:55 +00:00
start: {buildId: buildId, numberStr: ''},
2015-11-18 20:14:46 +00:00
}, this.slot());
},
function(err, lines) {
client.emit('sync', 'data', {lines: lines});
this.pass(true);
},
function(err) {
if (err) {
logger.error(
'error during read log for "' + buildId + '":',
err.stack || err
);
}
2015-11-18 20:14:46 +00:00
callback();
}
);
});
buildDataResourcesHash[buildId] = buildDataResource;
};
exports.createBuildDataResource = createBuildDataResource;
distributor.on('buildUpdate', function(build, changes) {
var buildsResource = app.dataio.resource('builds');
if (build.status === 'queued') {
createBuildDataResource(build.id);
}
2015-09-26 20:49:03 +00:00
// notify about build's project change, coz building affects project
// related stat (last build date, avg build time, etc)
if (changes.completed) {
var projectsResource = app.dataio.resource('projects');
projectsResource.clientEmitSyncChange({name: build.project.name});
}
buildsResource.clientEmitSync('change', {
buildId: build.id, changes: changes
});
});
2015-10-03 14:14:41 +00:00
var buildLogLineNumbersHash = {};
distributor.on('buildData', function(build, data) {
2015-11-21 19:33:41 +00:00
var lines = _(data.split('\n')).chain().invoke('trim').compact().value(),
2015-11-18 20:14:46 +00:00
logLineNumber = buildLogLineNumbersHash[build.id] || 0;
lines = _(lines).map(function(line, index) {
return {
number: logLineNumber + index,
text: line
};
});
buildLogLineNumbersHash[build.id] = logLineNumber + lines.length;
2015-10-15 05:36:04 +00:00
app.dataio.resource('build' + build.id).clientEmitSync(
2015-10-15 05:36:04 +00:00
'data',
2015-11-18 20:14:46 +00:00
{lines: lines}
);
2015-10-03 14:14:41 +00:00
2015-11-18 20:14:46 +00:00
_(lines).each(function(line) {
2015-11-19 17:20:35 +00:00
line.buildId = build.id;
});
// write build logs to db
2015-11-21 19:33:41 +00:00
if (lines.length) {
db.logLines.put(lines, function(err) {
if (err) {
logger.error(
'Error during write log line "' + logLineNumber +
'" for build "' + build.id + '":',
err.stack || err
);
}
});
}
});
callback(null, distributor);
};