Merge pull request #129 from ripple/tefalready-multiserver

Tefalready multiserver
This commit is contained in:
shekenahglory
2014-08-12 11:23:53 -07:00
3 changed files with 47 additions and 34 deletions

View File

@@ -788,6 +788,20 @@ Remote.prototype.isConnected = function() {
return this._connected;
};
/**
* Get array of connected servers
*/
Remote.prototype.getConnectedServers = function() {
var servers = [ ];
for (var i=0; i<this._servers.length; i++) {
if (this._servers[i].isConnected()) {
servers.push(this._servers[i]);
}
}
return servers;
};
/**
* Select a server to handle a request. Servers are
* automatically prioritized
@@ -795,41 +809,28 @@ Remote.prototype.isConnected = function() {
Remote.prototype._getServer =
Remote.prototype.getServer = function() {
var result = void(0);
if (this._primary_server && this._primary_server.isConnected()) {
return this._primary_server;
}
if (!this._servers.length) {
return result;
return null;
}
function sortByScore(a, b) {
var aScore = a._score + a._fee;
var bScore = b._score + b._fee;
if (aScore > bScore) {
return 1;
} else if (aScore < bScore) {
return -1;
} else {
return 0;
}
};
var connectedServers = this.getConnectedServers();
var server = connectedServers[0];
var cScore = server._score + server._fee;
// Sort servers by score
this._servers.sort(sortByScore);
// First connected server
for (var i=0; i<this._servers.length; i++) {
var server = this._servers[i];
if ((server instanceof Server) && server._connected) {
result = server;
break;
for (var i=1; i<connectedServers.length; i++) {
var _server = connectedServers[i];
var bScore = _server._score + _server._fee;
if (bScore < cScore) {
server = _server;
cScore = bScore;
}
}
return result;
return server;
};
/**
@@ -1479,15 +1480,15 @@ Remote.prototype.requestTxHistory = function(start, callback) {
*/
Remote.prototype.requestBookOffers = function(gets, pays, taker, callback) {
if (gets.hasOwnProperty('pays')) {
if (gets.hasOwnProperty('gets') || gets.hasOwnProperty('taker_gets')) {
var options = gets;
// This would mutate the `lastArg` in `arguments` to be `null` and is
// redundant. Once upon a time, some awkward code was written f(g, null,
// null, cb) ...
// callback = pays;
taker = options.taker;
pays = options.pays;
gets = options.gets;
pays = options.pays || options.taker_pays;
gets = options.gets || options.taker_gets;
}
var lastArg = arguments[arguments.length - 1];

View File

@@ -33,23 +33,27 @@ function Request(remote, command) {
util.inherits(Request, EventEmitter);
Request.prototype.broadcast = function() {
this._broadcast = true;
return this.request();
var connectedServers = this.remote.getConnectedServers();
this.request(connectedServers);
return connectedServers.length;
};
// Send the request to a remote.
Request.prototype.request = function(callback) {
Request.prototype.request = function(servers, callback) {
if (this.requested) {
return this;
}
this.requested = true;
if (typeof servers === 'function') {
callback = servers;
}
this.requested = true;
this.on('error', function(){});
this.emit('request', this.remote);
if (this._broadcast) {
this.remote._servers.forEach(function(server) {
if (Array.isArray(servers)) {
servers.forEach(function(server) {
this.setServer(server);
this.remote.request(this);
}, this);

View File

@@ -386,6 +386,11 @@ TransactionManager.prototype._request = function(tx) {
case 'tefPAST_SEQ':
self._resubmit(1, tx);
break;
case 'tefALREADY':
if (tx.responses === tx.submissions) {
tx.emit('error', message);
}
break;
default:
tx.emit('error', message);
}
@@ -450,6 +455,7 @@ TransactionManager.prototype._request = function(tx) {
message.result = message.engine_result || '';
tx.result = message;
tx.responses += 1;
if (remote.trace) {
log.info('submit response:', message);
@@ -543,7 +549,7 @@ TransactionManager.prototype._request = function(tx) {
}
submitRequest.timeout(self._submissionTimeout, requestTimeout);
submitRequest.broadcast();
tx.submissions = submitRequest.broadcast();
tx.attempts++;
tx.emit('postsubmit');
@@ -656,6 +662,8 @@ TransactionManager.prototype.submit = function(tx) {
}
tx.attempts = 0;
tx.submissions = 0;
tx.responses = 0;
// ND: this is the ONLY place we put the tx into the queue. The
// TransactionQueue queue is merely a list, so any mutations to tx._hash