better loader management.

This commit is contained in:
Christopher Jeffrey 2016-01-19 17:52:29 -08:00
parent 8abaef49a9
commit c11a651e9d
2 changed files with 129 additions and 43 deletions

View File

@ -24,11 +24,11 @@ try {
* Peer * Peer
*/ */
function Peer(pool, createSocket, options) { function Peer(pool, createConnection, options) {
var self = this; var self = this;
if (!(this instanceof Peer)) if (!(this instanceof Peer))
return new Peer(pool, createSocket, options); return new Peer(pool, createConnection, options);
EventEmitter.call(this); EventEmitter.call(this);
@ -42,6 +42,7 @@ function Peer(pool, createSocket, options) {
this.version = null; this.version = null;
this.destroyed = false; this.destroyed = false;
this.ack = false; this.ack = false;
this.connected = false;
this.ts = this.options.ts || 0; this.ts = this.options.ts || 0;
this.host = null; this.host = null;
@ -52,13 +53,13 @@ function Peer(pool, createSocket, options) {
if (this.options.backoff) { if (this.options.backoff) {
setTimeout(function() { setTimeout(function() {
self.socket = createSocket(self, pool); self.socket = createConnection(self, pool, options);
if (!self.socket) if (!self.socket)
throw new Error('No socket'); throw new Error('No socket');
self.emit('socket'); self.emit('socket');
}, this.options.backoff); }, this.options.backoff);
} else { } else {
this.socket = createSocket(this, pool); this.socket = createConnection(this, pool, options);
if (!this.socket) if (!this.socket)
throw new Error('No socket'); throw new Error('No socket');
} }
@ -107,10 +108,12 @@ Peer.prototype._init = function init() {
self.host = self.socket.remoteAddress; self.host = self.socket.remoteAddress;
if (!self.port) if (!self.port)
self.port = self.socket.remotePort; self.port = self.socket.remotePort;
self.emit('connect');
}); });
this.socket.once('error', function(err) { this.socket.once('error', function(err) {
self._error(err); self._error(err);
self.emit('misbehave');
}); });
this.socket.once('close', function() { this.socket.once('close', function() {
@ -131,6 +134,7 @@ Peer.prototype._init = function init() {
// Something is wrong here. // Something is wrong here.
// Ignore this peer. // Ignore this peer.
self.destroy(); self.destroy();
self.emit('misbehave');
}); });
if (this.pool.options.fullNode) { if (this.pool.options.fullNode) {
@ -152,12 +156,17 @@ Peer.prototype._init = function init() {
if (err) { if (err) {
self._error(err); self._error(err);
self.destroy(); self.destroy();
self.emit('misbehave');
return; return;
} }
self.ack = true; self.ack = true;
self.emit('ack'); self.emit('ack');
self.ts = utils.now(); self.ts = utils.now();
self._write(self.framer.packet('getaddr', [])); self._write(self.framer.packet('getaddr', []));
// if (self.pool.options.headers) {
// if (self.version.v > 70012)
// self._write(self.framer.packet('sendheaders', []));
// }
}); });
// Send hello // Send hello
@ -291,8 +300,13 @@ Peer.prototype._error = function error(err) {
if (this.destroyed) if (this.destroyed)
return; return;
if (typeof err === 'string')
err = new Error(err);
err.message += ' (' + this.host + ')';
this.destroy(); this.destroy();
this.emit('error', typeof err === 'string' ? new Error(err) : err); this.emit('error', err);
}; };
Peer.prototype._req = function _req(cmd, cb) { Peer.prototype._req = function _req(cmd, cb) {

View File

@ -37,9 +37,12 @@ function Pool(options) {
this.options.relay = this.options.relay == null this.options.relay = this.options.relay == null
? (this.options.fullNode ? true : false) ? (this.options.fullNode ? true : false)
: this.options.relay; : this.options.relay;
this.options.seeds = options.seeds || network.seeds;
this.setSeeds(this.options.seeds); this.originalSeeds = (options.seeds || network.seeds).map(utils.parseHost);
this.setSeeds([]);
this._priorityTries = {};
this._regularTries = {};
this._misbehaving = {};
this.storage = this.options.storage; this.storage = this.options.storage;
this.destroyed = false; this.destroyed = false;
@ -266,26 +269,32 @@ Pool.prototype._stopInterval = function _stopInterval() {
delete this._interval; delete this._interval;
}; };
Pool.prototype.createConnection = function createConnection(peer, pool) { Pool.prototype.createConnection = function createConnection(peer, pool, options) {
var addr, net, socket; var addr, net, socket;
if (pool._createConnection) if (pool._createConnection)
return pool._createConnection(peer, pool); return pool._createConnection(peer, pool);
addr = pool.usableSeed(true); addr = pool.usableSeed(options.priority, true);
peer.host = addr.host; peer.host = addr.host;
peer.port = addr.port; peer.port = addr.port;
if (this._createSocket) { if (pool._createSocket) {
socket = this._createSocket(addr.port, addr.host); socket = pool._createSocket(addr.port, addr.host);
} else { } else {
net = require('net'); net = require('net');
socket = net.connect(addr.port, addr.host); socket = net.connect(addr.port, addr.host);
} }
pool.emit('debug',
'Connecting to %s:%d (priority=%s)',
addr.host, addr.port, options.priority);
socket.on('connect', function() { socket.on('connect', function() {
pool.emit('debug', 'Connected to %s:%d', addr.host, addr.port); pool.emit('debug',
'Connected to %s:%d (priority=%s)',
addr.host, addr.port, options.priority);
}); });
return socket; return socket;
@ -301,7 +310,7 @@ Pool.prototype._addLoader = function _addLoader() {
if (this.peers.load != null) if (this.peers.load != null)
return; return;
peer = this._createPeer(750 * Math.random()); peer = this._createPeer(750 * Math.random(), true);
peer.once('socket', function() { peer.once('socket', function() {
self.emit('debug', 'Added loader peer: %s', peer.host); self.emit('debug', 'Added loader peer: %s', peer.host);
@ -658,13 +667,14 @@ Pool.prototype.loadMempool = function loadMempool() {
}); });
}; };
Pool.prototype._createPeer = function _createPeer(backoff) { Pool.prototype._createPeer = function _createPeer(backoff, priority) {
var self = this; var self = this;
var peer = new bcoin.peer(this, this.createConnection, { var peer = new bcoin.peer(this, this.createConnection, {
backoff: backoff, backoff: backoff,
startHeight: this.options.startHeight, startHeight: this.options.startHeight,
relay: this.options.relay relay: this.options.relay,
priority: priority
}); });
peer._retry = 0; peer._retry = 0;
@ -678,15 +688,19 @@ Pool.prototype._createPeer = function _createPeer(backoff) {
self.emit.apply(self, ['debug'].concat(args)); self.emit.apply(self, ['debug'].concat(args));
}); });
peer.once('misbehave', function() {
self._misbehaving[peer.host] = true;
});
peer.on('reject', function(payload) { peer.on('reject', function(payload) {
payload.data = utils.revHex(utils.toHex(payload.data)); var data = utils.revHex(utils.toHex(payload.data));
self.emit('debug', self.emit('debug',
'Reject: msg=%s ccode=%s reason=%s data=%s', 'Reject: msg=%s ccode=%s reason=%s data=%s',
payload.message, payload.message,
payload.ccode, payload.ccode,
payload.reason, payload.reason,
payload.data); data);
self.emit('reject', payload, peer); self.emit('reject', payload, peer);
}); });
@ -709,7 +723,7 @@ Pool.prototype._createPeer = function _createPeer(backoff) {
peer.on('addr', function(data) { peer.on('addr', function(data) {
if (self.seeds.length > 1000) if (self.seeds.length > 1000)
self.setSeeds(self.options.seeds.concat(self.seeds.slice(-500))); self.setSeeds(self.seeds.slice(-500));
self.addSeed(data); self.addSeed(data);
@ -724,6 +738,9 @@ Pool.prototype._createPeer = function _createPeer(backoff) {
if (version.height > self.block.bestHeight) if (version.height > self.block.bestHeight)
self.block.bestHeight = version.height; self.block.bestHeight = version.height;
self.emit('version', version, peer); self.emit('version', version, peer);
self.emit('debug',
'Received version from %s: version=%d height=%d agent=%s',
peer.host, peer.version.v, peer.version.height, peer.version.agent);
}); });
return peer; return peer;
@ -744,7 +761,7 @@ Pool.prototype._addPeer = function _addPeer(backoff) {
return; return;
} }
peer = this._createPeer(backoff); peer = this._createPeer(backoff, false);
this.peers.pending.push(peer); this.peers.pending.push(peer);
this.peers.all.push(peer); this.peers.all.push(peer);
@ -1477,40 +1494,95 @@ Pool.prototype.getPeer = function getPeer(addr) {
} }
}; };
Pool.prototype.usableSeed = function usableSeed(addrs, force) { Pool.prototype.usableSeed = function usableSeed(priority, connecting) {
var i, addr; var i, addr;
var original = this.originalSeeds;
var seeds = this.seeds;
var tries = priority ? this._priorityTries : this._regularTries;
if (typeof addrs === 'boolean') { // Hang back if we don't have a loader peer yet.
force = addrs; if (!connecting && !priority && (!this.peers.load || !this.peers.load.socket))
addrs = null; return;
// Randomize the non-original peers.
seeds = seeds.slice().sort(function() {
return Math.random() > 0.50 ? 1 : -1;
});
// Try to avoid connecting to a peer twice.
// Try the original peers first.
for (i = 0; i < original.length; i++) {
addr = original[i];
assert(addr.host);
if (this.getPeer(addr))
continue;
if (this._misbehaving[addr.host])
continue;
if (tries[addr.host])
continue;
if (connecting)
tries[addr.host] = true;
return addr;
} }
if (!addrs) // If we are a priority socket, try to find a
addrs = this.seeds; // peer this time with looser requirements.
if (priority) {
assert(addrs.length); for (i = 0; i < original.length; i++) {
addr = original[i];
if (this.peers.load) { assert(addr.host);
addrs = addrs.slice().sort(function() { if (this.peers.load && this.getPeer(addr) === this.peers.load)
return Math.random() > 0.50 ? 1 : -1; continue;
}); if (this._misbehaving[addr.host])
for (i = 0; i < addrs.length; i++) { continue;
if (!this.getPeer(addrs[i])) if (tries[addr.host])
return addrs[i]; continue;
if (connecting)
tries[addr.host] = true;
return addr;
} }
} }
if (addrs.length === 1) // Try the rest of the peers second.
return addrs[0]; for (i = 0; i < seeds.length; i++) {
addr = seeds[i];
assert(addr.host);
if (this.getPeer(addr))
continue;
if (this._misbehaving[addr.host])
continue;
return addr;
}
if (!force) // If we are a priority socket, try to find a
return; // peer this time with looser requirements.
if (priority) {
for (i = 0; i < seeds.length; i++) {
addr = seeds[i];
assert(addr.host);
if (this.peers.load && this.getPeer(addr) === this.peers.load)
continue;
if (this._misbehaving[addr.host])
continue;
return addr;
}
}
do { // If we have no block peers, always return
addr = addrs[Math.random() * (addrs.length - 1) | 0]; // an address.
} while (this.peers.load && this.getPeer(addr) === this.peers.load); if (!priority) {
if (this.peers.pending.length + this.peers.block.length === 0)
return original[Math.random() * (original.length - 1) | 0];
}
return addr; // This should never happen: priority sockets
// should _always_ get an address.
if (priority) {
this.emit('debug',
'We had to connect to a random peer. Something is not right.');
return seeds[Math.random() * (seeds.length - 1) | 0];
}
}; };
Pool.prototype.setSeeds = function setSeeds(seeds) { Pool.prototype.setSeeds = function setSeeds(seeds) {