2014-11-22 16:59:48 +01:00
|
|
|
'use strict';
|
|
|
|
|
2015-01-21 17:36:59 +01:00
|
|
|
const net = require('net');
|
|
|
|
const util = require('util');
|
2015-09-17 00:45:29 +02:00
|
|
|
const EventEmitter = require('events');
|
2015-01-21 17:36:59 +01:00
|
|
|
const debug = util.debuglog('http');
|
2013-04-11 23:47:15 +02:00
|
|
|
|
|
|
|
// New Agent code.
|
|
|
|
|
|
|
|
// The largest departure from the previous implementation is that
|
|
|
|
// an Agent instance holds connections for a variable number of host:ports.
|
|
|
|
// Surprisingly, this is still API compatible as far as third parties are
|
|
|
|
// concerned. The only code that really notices the difference is the
|
|
|
|
// request object.
|
|
|
|
|
|
|
|
// Another departure is that all code related to HTTP parsing is in
|
|
|
|
// ClientRequest.onSocket(). The Agent is now *strictly*
|
|
|
|
// concerned with managing a connection pool.
|
|
|
|
|
|
|
|
function Agent(options) {
|
2013-05-16 00:10:53 +02:00
|
|
|
if (!(this instanceof Agent))
|
|
|
|
return new Agent(options);
|
|
|
|
|
2013-04-11 23:47:15 +02:00
|
|
|
EventEmitter.call(this);
|
|
|
|
|
|
|
|
var self = this;
|
2013-05-23 03:44:24 +02:00
|
|
|
|
|
|
|
self.defaultPort = 80;
|
|
|
|
self.protocol = 'http:';
|
|
|
|
|
2013-05-21 23:02:18 +02:00
|
|
|
self.options = util._extend({}, options);
|
2013-05-23 03:44:24 +02:00
|
|
|
|
2013-05-21 23:02:18 +02:00
|
|
|
// don't confuse net and make it think that we're connecting to a pipe
|
|
|
|
self.options.path = null;
|
2013-04-11 23:47:15 +02:00
|
|
|
self.requests = {};
|
|
|
|
self.sockets = {};
|
2013-05-16 00:10:53 +02:00
|
|
|
self.freeSockets = {};
|
|
|
|
self.keepAliveMsecs = self.options.keepAliveMsecs || 1000;
|
|
|
|
self.keepAlive = self.options.keepAlive || false;
|
2013-04-11 23:47:15 +02:00
|
|
|
self.maxSockets = self.options.maxSockets || Agent.defaultMaxSockets;
|
2013-08-05 05:54:52 +02:00
|
|
|
self.maxFreeSockets = self.options.maxFreeSockets || 256;
|
2013-05-16 00:10:53 +02:00
|
|
|
|
2013-05-23 03:44:24 +02:00
|
|
|
self.on('free', function(socket, options) {
|
|
|
|
var name = self.getName(options);
|
|
|
|
debug('agent.on(free)', name);
|
2013-04-11 23:47:15 +02:00
|
|
|
|
2016-01-12 05:41:01 +01:00
|
|
|
if (socket.writable &&
|
2013-04-11 23:47:15 +02:00
|
|
|
self.requests[name] && self.requests[name].length) {
|
|
|
|
self.requests[name].shift().onSocket(socket);
|
|
|
|
if (self.requests[name].length === 0) {
|
|
|
|
// don't leak
|
|
|
|
delete self.requests[name];
|
|
|
|
}
|
|
|
|
} else {
|
2013-05-16 00:10:53 +02:00
|
|
|
// If there are no pending requests, then put it in
|
|
|
|
// the freeSockets pool, but only if we're allowed to do so.
|
|
|
|
var req = socket._httpMessage;
|
|
|
|
if (req &&
|
|
|
|
req.shouldKeepAlive &&
|
2016-01-12 05:41:01 +01:00
|
|
|
socket.writable &&
|
2015-12-24 00:52:01 +01:00
|
|
|
self.keepAlive) {
|
2013-05-16 00:10:53 +02:00
|
|
|
var freeSockets = self.freeSockets[name];
|
2013-08-05 05:54:52 +02:00
|
|
|
var freeLen = freeSockets ? freeSockets.length : 0;
|
|
|
|
var count = freeLen;
|
2013-05-16 00:10:53 +02:00
|
|
|
if (self.sockets[name])
|
|
|
|
count += self.sockets[name].length;
|
|
|
|
|
2015-03-23 12:33:13 +01:00
|
|
|
if (count > self.maxSockets || freeLen >= self.maxFreeSockets) {
|
2013-05-16 00:10:53 +02:00
|
|
|
socket.destroy();
|
|
|
|
} else {
|
|
|
|
freeSockets = freeSockets || [];
|
|
|
|
self.freeSockets[name] = freeSockets;
|
|
|
|
socket.setKeepAlive(true, self.keepAliveMsecs);
|
|
|
|
socket.unref();
|
|
|
|
socket._httpMessage = null;
|
2013-08-05 05:54:52 +02:00
|
|
|
self.removeSocket(socket, options);
|
2013-05-16 00:10:53 +02:00
|
|
|
freeSockets.push(socket);
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
socket.destroy();
|
|
|
|
}
|
2013-04-11 23:47:15 +02:00
|
|
|
}
|
|
|
|
});
|
|
|
|
}
|
2013-05-15 23:24:08 +02:00
|
|
|
|
2013-04-11 23:47:15 +02:00
|
|
|
util.inherits(Agent, EventEmitter);
|
|
|
|
exports.Agent = Agent;
|
|
|
|
|
2013-05-16 00:10:53 +02:00
|
|
|
Agent.defaultMaxSockets = Infinity;
|
2013-04-11 23:47:15 +02:00
|
|
|
|
2013-05-21 23:02:18 +02:00
|
|
|
Agent.prototype.createConnection = net.createConnection;
|
2013-05-23 03:44:24 +02:00
|
|
|
|
|
|
|
// Get the key for a given set of request options
|
2016-10-12 13:03:13 +02:00
|
|
|
Agent.prototype.getName = function getName(options) {
|
2015-09-11 23:18:01 +02:00
|
|
|
var name = options.host || 'localhost';
|
2013-05-23 03:44:24 +02:00
|
|
|
|
|
|
|
name += ':';
|
|
|
|
if (options.port)
|
|
|
|
name += options.port;
|
2015-09-11 23:18:01 +02:00
|
|
|
|
2013-05-23 03:44:24 +02:00
|
|
|
name += ':';
|
|
|
|
if (options.localAddress)
|
|
|
|
name += options.localAddress;
|
2015-09-11 23:18:01 +02:00
|
|
|
|
2016-05-09 16:35:20 +02:00
|
|
|
// Pacify parallel/test-http-agent-getname by only appending
|
|
|
|
// the ':' when options.family is set.
|
|
|
|
if (options.family === 4 || options.family === 6)
|
|
|
|
name += ':' + options.family;
|
|
|
|
|
2013-05-23 03:44:24 +02:00
|
|
|
return name;
|
|
|
|
};
|
|
|
|
|
2016-10-12 13:03:13 +02:00
|
|
|
Agent.prototype.addRequest = function addRequest(req, options) {
|
2016-02-11 01:59:25 +01:00
|
|
|
// Legacy API: addRequest(req, host, port, localAddress)
|
2013-08-07 19:23:45 +02:00
|
|
|
if (typeof options === 'string') {
|
|
|
|
options = {
|
|
|
|
host: options,
|
|
|
|
port: arguments[2],
|
2016-02-11 01:59:25 +01:00
|
|
|
localAddress: arguments[3]
|
2013-08-07 19:23:45 +02:00
|
|
|
};
|
|
|
|
}
|
|
|
|
|
2016-03-14 16:56:02 +01:00
|
|
|
options = util._extend({}, options);
|
2017-02-14 02:58:56 +01:00
|
|
|
util._extend(options, this.options);
|
2016-03-14 16:56:02 +01:00
|
|
|
|
2016-09-18 13:34:48 +02:00
|
|
|
if (!options.servername) {
|
|
|
|
options.servername = options.host;
|
|
|
|
const hostHeader = req.getHeader('host');
|
|
|
|
if (hostHeader) {
|
|
|
|
options.servername = hostHeader.replace(/:.*$/, '');
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2013-05-23 03:44:24 +02:00
|
|
|
var name = this.getName(options);
|
2013-04-11 23:47:15 +02:00
|
|
|
if (!this.sockets[name]) {
|
|
|
|
this.sockets[name] = [];
|
|
|
|
}
|
2013-05-16 00:10:53 +02:00
|
|
|
|
2013-08-05 05:54:52 +02:00
|
|
|
var freeLen = this.freeSockets[name] ? this.freeSockets[name].length : 0;
|
|
|
|
var sockLen = freeLen + this.sockets[name].length;
|
|
|
|
|
|
|
|
if (freeLen) {
|
2013-05-16 00:10:53 +02:00
|
|
|
// we have a free socket, so use that.
|
|
|
|
var socket = this.freeSockets[name].shift();
|
2013-08-05 05:54:52 +02:00
|
|
|
debug('have free socket');
|
2013-05-16 00:10:53 +02:00
|
|
|
|
|
|
|
// don't leak
|
|
|
|
if (!this.freeSockets[name].length)
|
|
|
|
delete this.freeSockets[name];
|
|
|
|
|
|
|
|
socket.ref();
|
|
|
|
req.onSocket(socket);
|
2013-08-05 05:54:52 +02:00
|
|
|
this.sockets[name].push(socket);
|
|
|
|
} else if (sockLen < this.maxSockets) {
|
|
|
|
debug('call onSocket', sockLen, freeLen);
|
2013-04-11 23:47:15 +02:00
|
|
|
// If we are under maxSockets create a new one.
|
2016-01-12 05:41:01 +01:00
|
|
|
this.createSocket(req, options, function(err, newSocket) {
|
|
|
|
if (err) {
|
|
|
|
process.nextTick(function() {
|
|
|
|
req.emit('error', err);
|
|
|
|
});
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
req.onSocket(newSocket);
|
|
|
|
});
|
2013-04-11 23:47:15 +02:00
|
|
|
} else {
|
2013-05-23 03:44:24 +02:00
|
|
|
debug('wait for socket');
|
2013-04-11 23:47:15 +02:00
|
|
|
// We are over limit so we'll add it to the queue.
|
|
|
|
if (!this.requests[name]) {
|
|
|
|
this.requests[name] = [];
|
|
|
|
}
|
|
|
|
this.requests[name].push(req);
|
|
|
|
}
|
|
|
|
};
|
2013-05-15 23:24:08 +02:00
|
|
|
|
2016-10-12 13:03:13 +02:00
|
|
|
Agent.prototype.createSocket = function createSocket(req, options, cb) {
|
2013-04-11 23:47:15 +02:00
|
|
|
var self = this;
|
2013-05-23 03:44:24 +02:00
|
|
|
options = util._extend({}, options);
|
2017-02-14 02:58:56 +01:00
|
|
|
util._extend(options, self.options);
|
2013-04-11 23:47:15 +02:00
|
|
|
|
2015-03-09 21:00:24 +01:00
|
|
|
if (!options.servername) {
|
|
|
|
options.servername = options.host;
|
2016-01-12 05:41:01 +01:00
|
|
|
const hostHeader = req.getHeader('host');
|
|
|
|
if (hostHeader) {
|
|
|
|
options.servername = hostHeader.replace(/:.*$/, '');
|
2013-04-11 23:47:15 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2013-05-23 03:44:24 +02:00
|
|
|
var name = self.getName(options);
|
2015-07-23 06:18:38 +02:00
|
|
|
options._agentKey = name;
|
2013-05-23 03:44:24 +02:00
|
|
|
|
|
|
|
debug('createConnection', name, options);
|
2013-11-03 23:32:03 +01:00
|
|
|
options.encoding = null;
|
2016-01-12 05:41:01 +01:00
|
|
|
var called = false;
|
|
|
|
const newSocket = self.createConnection(options, oncreate);
|
|
|
|
if (newSocket)
|
|
|
|
oncreate(null, newSocket);
|
|
|
|
function oncreate(err, s) {
|
|
|
|
if (called)
|
|
|
|
return;
|
|
|
|
called = true;
|
|
|
|
if (err)
|
|
|
|
return cb(err);
|
|
|
|
if (!self.sockets[name]) {
|
|
|
|
self.sockets[name] = [];
|
|
|
|
}
|
|
|
|
self.sockets[name].push(s);
|
|
|
|
debug('sockets', name, self.sockets[name].length);
|
2016-12-05 02:51:51 +01:00
|
|
|
installListeners(self, s, options);
|
2016-01-12 05:41:01 +01:00
|
|
|
cb(null, s);
|
2013-04-11 23:47:15 +02:00
|
|
|
}
|
|
|
|
};
|
2013-05-15 23:24:08 +02:00
|
|
|
|
2016-12-05 02:51:51 +01:00
|
|
|
function installListeners(agent, s, options) {
|
|
|
|
function onFree() {
|
|
|
|
debug('CLIENT socket onFree');
|
|
|
|
agent.emit('free', s, options);
|
|
|
|
}
|
|
|
|
s.on('free', onFree);
|
|
|
|
|
|
|
|
function onClose(err) {
|
|
|
|
debug('CLIENT socket onClose');
|
|
|
|
// This is the only place where sockets get removed from the Agent.
|
|
|
|
// If you want to remove a socket from the pool, just close it.
|
|
|
|
// All socket errors end in a close event anyway.
|
|
|
|
agent.removeSocket(s, options);
|
|
|
|
}
|
|
|
|
s.on('close', onClose);
|
|
|
|
|
|
|
|
function onRemove() {
|
|
|
|
// We need this function for cases like HTTP 'upgrade'
|
|
|
|
// (defined by WebSockets) where we need to remove a socket from the
|
|
|
|
// pool because it'll be locked up indefinitely
|
|
|
|
debug('CLIENT socket onRemove');
|
|
|
|
agent.removeSocket(s, options);
|
|
|
|
s.removeListener('close', onClose);
|
|
|
|
s.removeListener('free', onFree);
|
|
|
|
s.removeListener('agentRemove', onRemove);
|
|
|
|
}
|
|
|
|
s.on('agentRemove', onRemove);
|
|
|
|
}
|
|
|
|
|
2016-10-12 13:03:13 +02:00
|
|
|
Agent.prototype.removeSocket = function removeSocket(s, options) {
|
2013-05-23 03:44:24 +02:00
|
|
|
var name = this.getName(options);
|
2016-01-12 05:41:01 +01:00
|
|
|
debug('removeSocket', name, 'writable:', s.writable);
|
2013-11-06 12:58:03 +01:00
|
|
|
var sets = [this.sockets];
|
|
|
|
|
|
|
|
// If the socket was destroyed, remove it from the free buffers too.
|
2016-01-12 05:41:01 +01:00
|
|
|
if (!s.writable)
|
2013-11-06 12:58:03 +01:00
|
|
|
sets.push(this.freeSockets);
|
|
|
|
|
2014-08-14 05:15:24 +02:00
|
|
|
for (var sk = 0; sk < sets.length; sk++) {
|
|
|
|
var sockets = sets[sk];
|
|
|
|
|
2013-11-06 12:58:03 +01:00
|
|
|
if (sockets[name]) {
|
|
|
|
var index = sockets[name].indexOf(s);
|
|
|
|
if (index !== -1) {
|
|
|
|
sockets[name].splice(index, 1);
|
|
|
|
// Don't leak
|
|
|
|
if (sockets[name].length === 0)
|
|
|
|
delete sockets[name];
|
2013-04-11 23:47:15 +02:00
|
|
|
}
|
|
|
|
}
|
2014-08-14 05:15:24 +02:00
|
|
|
}
|
|
|
|
|
2013-04-11 23:47:15 +02:00
|
|
|
if (this.requests[name] && this.requests[name].length) {
|
2013-05-23 03:44:24 +02:00
|
|
|
debug('removeSocket, have a request, make a socket');
|
2013-04-11 23:47:15 +02:00
|
|
|
var req = this.requests[name][0];
|
2013-05-23 03:44:24 +02:00
|
|
|
// If we have pending requests and a socket gets closed make a new one
|
2016-01-12 05:41:01 +01:00
|
|
|
this.createSocket(req, options, function(err, newSocket) {
|
|
|
|
if (err) {
|
|
|
|
process.nextTick(function() {
|
|
|
|
req.emit('error', err);
|
|
|
|
});
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
newSocket.emit('free');
|
|
|
|
});
|
2013-04-11 23:47:15 +02:00
|
|
|
}
|
|
|
|
};
|
|
|
|
|
2016-10-12 13:03:13 +02:00
|
|
|
Agent.prototype.destroy = function destroy() {
|
2013-05-16 00:10:53 +02:00
|
|
|
var sets = [this.freeSockets, this.sockets];
|
2014-08-14 05:15:24 +02:00
|
|
|
for (var s = 0; s < sets.length; s++) {
|
|
|
|
var set = sets[s];
|
|
|
|
var keys = Object.keys(set);
|
|
|
|
for (var v = 0; v < keys.length; v++) {
|
|
|
|
var setName = set[keys[v]];
|
|
|
|
for (var n = 0; n < setName.length; n++) {
|
|
|
|
setName[n].destroy();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2013-05-16 00:10:53 +02:00
|
|
|
};
|
|
|
|
|
2013-08-05 21:33:19 +02:00
|
|
|
exports.globalAgent = new Agent();
|