From 46db9d9d7b668a8d7bd3aa69390c7734a98e4151 Mon Sep 17 00:00:00 2001 From: Jonathan Gros-Dubois Date: Mon, 23 Sep 2013 05:11:47 -0700 Subject: [PATCH 0001/1219] Initial commit --- README.md | 4 ++++ 1 file changed, 4 insertions(+) create mode 100644 README.md diff --git a/README.md b/README.md new file mode 100644 index 00000000..b6f2e9e5 --- /dev/null +++ b/README.md @@ -0,0 +1,4 @@ +socketcluster +============= + +Highly scalable realtime sockets based on Engine.io From e4727a3183570e4299c2762a73c2c1012f3a8bee Mon Sep 17 00:00:00 2001 From: Jonathan Gros-Dubois Date: Mon, 23 Sep 2013 22:14:34 +1000 Subject: [PATCH 0002/1219] First commit --- README.md | 170 ++++++++++- balancer.js | 69 +++++ index.js | 525 ++++++++++++++++++++++++++++++++++ node_modules/.bin/express | 15 + node_modules/.bin/express.cmd | 5 + node_modules/.gitignore | 4 + node_modules/json/cycle.js | 186 ++++++++++++ node_modules/json/index.js | 493 +++++++++++++++++++++++++++++++ package.json | 30 ++ socketworker.js | 205 +++++++++++++ worker-bootstrap.js | 55 ++++ 11 files changed, 1754 insertions(+), 3 deletions(-) create mode 100644 balancer.js create mode 100644 index.js create mode 100644 node_modules/.bin/express create mode 100644 node_modules/.bin/express.cmd create mode 100644 node_modules/.gitignore create mode 100644 node_modules/json/cycle.js create mode 100644 node_modules/json/index.js create mode 100644 package.json create mode 100644 socketworker.js create mode 100644 worker-bootstrap.js diff --git a/README.md b/README.md index b6f2e9e5..9120119b 100644 --- a/README.md +++ b/README.md @@ -1,4 +1,168 @@ -socketcluster -============= +SocketCluster +====== -Highly scalable realtime sockets based on Engine.io +SocketCluster is a WebSocket server cluster (with HTTP long-polling fallback) based on engine.io. +Unlike other realtime engines, SocketCluster deploys itself as cluster in order to make use of all CPUs/cores on +a machine/instance - This offers a more consistent performance for users and lets you scale vertically without limits. +SocketCluster workers are highly parallelized - Asymptotically speaking, SocketCluster is N times faster than any other +available WebSocket server (where N is the number of CPUs/cores available on your machine). +SocketCluster was designed to be lightweight and its API is almost identical to Socket.io. + +Other advantages of SocketCluster include: +- Sockets which are bound to the same browser (for example, across multiple tabs) share the same session. +- You can emit an event on a session to notify all sockets that belong to it. +- The SocketCluster client (socketcluster-client) has an option to allow disconnected sockets to automatically (and seamlessly) reconnect +if they lose the connection. +- Server crashes are transparent to users (aside from a 2 to 5 second delay to allow the worker to respawn) - Session data remains intact between crashes. +- It uses a memory store cluster called nData which you can use to store 'volatile' session data which relates to your sockets/sessions. + +To install, run: + +```bash +npm install socketcluster +``` + +## How to use + +The following example launches SocketCluster as 7 distinct processes (in addition to the current master process): +- 3 workers on ports 9100, 9101, 9102 +- 3 stores on ports 9001, 9002, 9003 +- 1 load balancer on port 8000 which distributes requests evenly between the 3 workers + +```js +var socketCluster = new SocketCluster({ + workers: [9100, 9101, 9102], + stores: [9001, 9002, 9003], + balancerCount: 1, // Optional + port: 8000, + appName: 'myapp', + workerController: 'worker.js' +}); +``` + +The appName option can be any string which uniquely identifies this application. +This avoids potential issues with having multiple SocketCluster apps run under the same domain +- It is used internally for various purposes. + +The workerController option is the path to a file which each SocketCluster worker will use to bootstrap itself. +This file is a standard Node.js module which must expose a run() function - Inside this run function is where you should +put all your application logic. + +Example 'worker.js': + +```js +var fs = require('fs'); + +module.exports.run = function (worker) { + // Get a reference to our raw Node HTTP server + var httpServer = worker.getHTTPServer(); + // Get a reference to our WebSocket server + var wsServer = worker.getWSServer(); + + /* + We're going to read our main HTML file and the socketcluster-client + script from disk and serve it to clients using the Node HTTP server. + */ + + var htmlPath = __dirname + '/index.html'; + var clientPath = __dirname + '/node_modules/socketcluster-client/socketcluster.js'; + + var html = fs.readFileSync(htmlPath, { + encoding: 'utf8' + }); + + var clientCode = fs.readFileSync(clientPath, { + encoding: 'utf8' + }); + + /* + Very basic code to serve our main HTML file to clients and + our socketcluster-client script when requested. + It may be better to use a framework like express here. + Note that the 'req' event used here is different from the standard Node.js HTTP server 'request' event + - The 'request' event also captures SocketCluster-related requests; the 'req' + event only captures the ones you actually need. As a rule of thumb, you should not listen to the 'request' event. + */ + httpServer.on('req', function (req, res) { + if (req.url == '/socketcluster.js') { + res.writeHead(200, { + 'Content-Type': 'text/javascript' + }); + res.end(clientCode); + } else if (req.url == '/') { + res.writeHead(200, { + 'Content-Type': 'text/html' + }); + res.end(html); + } + }); + + /* + In here we handle our incoming WebSocket connections and listen for events. + From here onwards is just like Socket.io. + */ + wsServer.on('connection', function (socket) { + setInterval(function () { + socket.emit('rand', Math.floor(Math.random() * 100)); + }, 1000); + }); +}; +``` + +### Using with Express + +Using SocketCluster with express is simple, you put the code inside your workerController: + +```js +module.exports.run = function (worker) { + // Get a reference to our raw Node HTTP server + var httpServer = worker.getHTTPServer(); + // Get a reference to our WebSocket server + var wsServer = worker.getWSServer(); + + var app = require('express')(); + + // Add whatever express middleware you like... + + // Make your express app handle all essential requests + httpServer.on('req', app); +} +``` + +## API (Coming soon) + +### SocketCluster + + Exposed by `require('socketcluster').SocketCluster`. + +### SocketCluster(opts:Object) + Creates a new SocketCluster, must be invoked with the new keyword. + + ```js + var SocketCluster = require('socketcluster').SocketCluster; + + var socketCluster = new SocketCluster({ + workers: [9100, 9101, 9102], + stores: [9001, 9002, 9003], + port: 8000, + appName: 'myapp', + workerController: 'worker.js' + }); + ``` + + Documentation on all supported options is coming soon (there are around 30 of them - Most of them are optional). + +### SocketWorker + + A SocketWorker object is passed as the argument to your workerController's run(worker) function. + Example - Inside worker.js: + + ```js + module.exports.run = function (worker) { + // worker here is an instance of SocketWorker + } + ``` + +### ClusterServer + + A ClusterServer instance is returned from worker.getWSServer() - You use it to handle WebSocket connections. \ No newline at end of file diff --git a/balancer.js b/balancer.js new file mode 100644 index 00000000..06df80d6 --- /dev/null +++ b/balancer.js @@ -0,0 +1,69 @@ +var cluster = require('cluster'); +var LoadBalancer = require('loadbalancer'); + +var balancer; + +if (cluster.isMaster) { + process.on('message', function (m) { + var balancers; + if (m.type == 'init') { + var balancerCount = m.data.balancerCount; + balancers = []; + + var launchBalancer = function (i) { + balancer = cluster.fork(); + balancers[i] = balancer; + balancer.on('error', function (err) { + process.send({ + message: err.message, + stack: err.stack + }); + }); + + balancer.on('message', process.send.bind(process)); + balancer.on('exit', function () { + launchBalancer(i); + }) + balancer.send(m); + }; + + for (var i=0; i 3) { + process.env.DEBUG = 'engine*'; + } else if (self.options.logLevel > 2) { + process.env.DEBUG = 'engine'; + } + + if (self.options.protocolOptions) { + var protoOpts = self.options.protocolOptions; + if (protoOpts.key instanceof Buffer) { + protoOpts.key = protoOpts.key.toString(); + } + if (protoOpts.cert instanceof Buffer) { + protoOpts.cert = protoOpts.cert.toString(); + } + if (protoOpts.pfx instanceof Buffer) { + protoOpts.pfx = protoOpts.pfx.toString(); + } + if (protoOpts.passphrase == null) { + var privKeyEncLine = protoOpts.key.split('\n')[1]; + if (privKeyEncLine.toUpperCase().indexOf('ENCRYPTED') > -1) { + var message = 'The supplied private key is encrypted and cannot be used without a passphrase - ' + + 'Please provide a valid passphrase as a property to protocolOptions'; + throw new Error(message); + } + process.exit(); + } + } + + if (self.options.stores) { + var newStores = []; + var curStore; + + for (i in self.options.stores) { + curStore = self.options.stores[i]; + if (typeof curStore == 'number') { + curStore = {port: curStore}; + } else { + if (curStore.port == null) { + throw new Error('One or more store objects is missing a port property'); + } + } + newStores.push(curStore); + } + self.options.stores = newStores; + } + + if (!self.options.stores || self.options.stores.length < 1) { + self.options.stores = [{port: self.options.port + 2}]; + } + + if (self.options.workers) { + var newWorkers = []; + var curWorker; + + for (i in self.options.workers) { + curWorker = self.options.workers[i]; + if (typeof curWorker == 'number') { + curWorker = {port: curWorker}; + } else { + if (curWorker.port == null) { + throw new Error('One or more worker objects is missing a port property'); + } + } + newWorkers.push(curWorker); + } + self.options.workers = newWorkers; + } else { + self.options.workers = [{port: self.options.port + 3}]; + } + + if (!self.options.balancerCount) { + self.options.balancerCount = Math.floor(self.options.workers.length / 2); + if (self.options.balancerCount < 1) { + self.options.balancerCount = 1; + } + } + + self._extRegex = /[.][^\/\\]*$/; + self._slashSequenceRegex = /\/+/g; + self._startSlashRegex = /^\//; + + self._minAddressSocketLimit = 20; + self._dataExpiryAccuracy = 5000; + + if (self.options.addressSocketLimit == null) { + var limit = self.options.sessionTimeout / 40; + if (limit < self._minAddressSocketLimit) { + limit = self._minAddressSocketLimit; + } + self.options.addressSocketLimit = limit; + } + + self._clusterEngine = require(self.options.clusterEngine); + + self._colorCodes = { + red: 31, + green: 32, + yellow: 33 + }; + + if (self.options.logLevel > 0) { + console.log(' ' + self.colorText('[Busy]', 'yellow') + ' Launching SocketCluster'); + } + + process.stdin.on('error', function (err) { + self.noticeHandler(err, {type: 'master'}); + }); + process.stdin.resume(); + process.stdin.setEncoding('utf8'); + + self._start(); +}; + +SocketCluster.prototype.errorHandler = function (err, origin) { + if (err.stack == null) { + if (!(err instanceof Object)) { + err = new Error(err); + } + err.stack = err.message; + } + if (origin instanceof Object) { + err.origin = origin; + } else { + err.origin = { + type: origin + }; + } + err.time = Date.now(); + + this.emit(this.EVENT_FAIL, err); + this.log(err.stack); +}; + +SocketCluster.prototype.noticeHandler = function (notice, origin) { + if (notice.stack == null) { + if (!(notice instanceof Object)) { + notice = new Error(notice); + } + notice.stack = notice.message; + } + if (origin instanceof Object) { + notice.origin = origin; + } else { + notice.origin = { + type: origin + }; + } + notice.time = Date.now(); + + this.emit(this.EVENT_NOTICE, notice); +}; + +SocketCluster.prototype.triggerInfo = function (info, origin) { + if (this._active) { + if (!(origin instanceof Object)) { + origin = { + type: origin + }; + } + var infoData = { + origin: origin, + info: info, + time: Date.now() + }; + this.emit(this.EVENT_INFO, infoData); + + if (this.options.logLevel > 0) { + this.log(info, infoData.time); + } + } +}; + +SocketCluster.prototype._start = function () { + var self = this; + + self._workers = []; + self._active = false; + + var leaderId = -1; + var firstTime = true; + + var workersActive = false; + + var initLoadBalancer = function () { + self._balancer.send({ + type: 'init', + data: { + dataKey: pass, + sourcePort: self.options.port, + workers: self.options.workers, + host: self.options.host, + balancerCount: self.options.balancerCount, + protocol: self.options.protocol, + protocolOptions: self.options.protocolOptions, + checkStatusTimeout: self.options.connectTimeout * 1000, + statusURL: self._paths.statusURL, + statusCheckInterval: self.options.workerStatusInterval * 1000 + } + }); + }; + + var launchLoadBalancer = function () { + if (self._balancer) { + self._errorDomain.remove(self._balancer); + } + + var balancerErrorHandler = function (err) { + self.errorHandler(err, {type: 'balancer'}); + }; + + var balancerNoticeHandler = function (noticeMessage) { + self.noticeHandler(noticeMessage, {type: 'balancer'}); + }; + + self._balancer = fork(__dirname + '/balancer.js'); + self._balancer.on('error', balancerErrorHandler); + self._balancer.on('notice', balancerNoticeHandler); + + self._balancer.on('exit', launchLoadBalancer); + self._balancer.on('message', function (m) { + if (m.type == 'error') { + balancerErrorHandler(m.data); + } else if (m.type == 'notice') { + balancerNoticeHandler(m.data); + } + }); + + if (workersActive) { + initLoadBalancer(); + } + }; + + launchLoadBalancer(); + + var workerIdCounter = 1; + + var ioClusterReady = function () { + var i; + var workerReadyHandler = function (data, worker) { + self._workers.push(worker); + if (worker.id == leaderId) { + worker.send({ + type: 'emit', + event: self.EVENT_LEADER_START + }); + } + + if (self._active && self.options.logLevel > 0) { + self.log('Worker ' + worker.data.id + ' was respawned on port ' + worker.data.port); + } + + if (self._workers.length >= self.options.workers.length) { + if (firstTime) { + if (self.options.logLevel > 0) { + console.log(' ' + self.colorText('[Active]', 'green') + ' SocketCluster started'); + console.log(' Port: ' + self.options.port); + console.log(' Master PID: ' + process.pid); + console.log(' Balancer count: ' + self.options.balancerCount); + console.log(' Worker count: ' + self.options.workers.length); + console.log(' Store count: ' + self.options.stores.length); + console.log(); + } + firstTime = false; + + if (!workersActive) { + initLoadBalancer(); + workersActive = true; + } + + process.on('SIGUSR2', function () { + for (var i in self._workers) { + self._workers[i].kill(); + } + }); + } else { + var workersData = []; + var i; + for (i in self._workers) { + workersData.push(self._workers[i].data); + } + self._balancer.send({ + type: 'setWorkers', + data: workersData + }); + } + self._active = true; + } + }; + + var launchWorker = function (workerData, lead) { + var workerErrorHandler = function (err) { + var origin = { + type: 'worker', + pid: worker.pid + }; + self.errorHandler(err, origin); + }; + + var workerNoticeHandler = function (noticeMessage) { + var origin = { + type: 'worker', + pid: worker.pid + }; + self.noticeHandler(noticeMessage, origin); + }; + + var worker = fork(__dirname + '/worker-bootstrap.js'); + worker.on('error', workerErrorHandler); + + if (!workerData.id) { + workerData.id = workerIdCounter++; + } + + worker.id = workerData.id; + worker.data = workerData; + + var workerOpts = self._cloneObject(self.options); + workerOpts.paths = self._paths; + workerOpts.workerId = worker.id; + workerOpts.sourcePort = self.options.port; + workerOpts.workerPort = workerData.port; + workerOpts.stores = stores; + workerOpts.dataKey = pass; + workerOpts.lead = lead ? 1 : 0; + + worker.send({ + type: 'init', + data: workerOpts + }); + + worker.on('message', function workerHandler(m) { + if (m.type == 'ready') { + if (lead) { + leaderId = worker.id; + } + workerReadyHandler(m, worker); + } else if (m.type == 'error') { + workerErrorHandler(m.data); + } else if (m.type == 'notice') { + workerNoticeHandler(m.data); + } + }); + + worker.on('exit', function (code, signal) { + self._errorDomain.remove(worker); + var message = ' Worker ' + worker.id + ' died - Exit code: ' + code; + + if (signal) { + message += ', signal: ' + signal; + } + + var workersData = []; + var newWorkers = []; + var i; + for (i in self._workers) { + if (self._workers[i].id != worker.id) { + newWorkers.push(self._workers[i]); + workersData.push(self._workers[i].data); + } + } + + self._workers = newWorkers; + self._balancer.send({ + type: 'setWorkers', + data: workersData + }); + + var lead = worker.id == leaderId; + leaderId = -1; + self.errorHandler(new Error(message), {type: 'master'}); + + if (self.options.logLevel > 0) { + self.log('Respawning worker ' + worker.id); + } + launchWorker(workerData, lead); + }); + + return worker; + }; + + var len = self.options.workers.length; + if (len > 0) { + launchWorker(self.options.workers[0], true); + for (i = 1; i < len; i++) { + launchWorker(self.options.workers[i]); + } + } + }; + + var stores = self.options.stores; + var pass = crypto.randomBytes(32).toString('hex'); + + var launchIOCluster = function () { + self._ioCluster = new self._clusterEngine.IOCluster({ + stores: stores, + dataKey: pass, + expiryAccuracy: self._dataExpiryAccuracy + }); + + self._ioCluster.on('error', function (err) { + self.errorHandler(err, {type: 'store'}); + }); + }; + + launchIOCluster(); + self._ioCluster.on('ready', ioClusterReady); +}; + +SocketCluster.prototype.log = function (message, time) { + if (time == null) { + time = Date.now(); + } + console.log(time + ' - ' + message); +}; + +SocketCluster.prototype._cloneObject = function (object) { + var clone = {}; + for (var i in object) { + clone[i] = object[i]; + } + return clone; +}; + +SocketCluster.prototype.colorText = function (message, color) { + if (this._colorCodes[color]) { + return '\033[0;' + this._colorCodes[color] + 'm' + message + '\033[0m'; + } else if (color) { + return '\033[' + color + 'm' + message + '\033[0m'; + } + return message; +}; + +module.exports.SocketCluster = SocketCluster; \ No newline at end of file diff --git a/node_modules/.bin/express b/node_modules/.bin/express new file mode 100644 index 00000000..cad5a1ef --- /dev/null +++ b/node_modules/.bin/express @@ -0,0 +1,15 @@ +#!/bin/sh +basedir=`dirname "$0"` + +case `uname` in + *CYGWIN*) basedir=`cygpath -w "$basedir"`;; +esac + +if [ -x "$basedir/node" ]; then + "$basedir/node" "$basedir/../express/bin/express" "$@" + ret=$? +else + node "$basedir/../express/bin/express" "$@" + ret=$? +fi +exit $ret diff --git a/node_modules/.bin/express.cmd b/node_modules/.bin/express.cmd new file mode 100644 index 00000000..1d2f5091 --- /dev/null +++ b/node_modules/.bin/express.cmd @@ -0,0 +1,5 @@ +@IF EXIST "%~dp0\node.exe" ( + "%~dp0\node.exe" "%~dp0\..\express\bin\express" %* +) ELSE ( + node "%~dp0\..\express\bin\express" %* +) \ No newline at end of file diff --git a/node_modules/.gitignore b/node_modules/.gitignore new file mode 100644 index 00000000..4d4c483b --- /dev/null +++ b/node_modules/.gitignore @@ -0,0 +1,4 @@ +iocluster/ +loadbalancer/ +ndata/ +socketcluster-server/ \ No newline at end of file diff --git a/node_modules/json/cycle.js b/node_modules/json/cycle.js new file mode 100644 index 00000000..44bfab09 --- /dev/null +++ b/node_modules/json/cycle.js @@ -0,0 +1,186 @@ +/* + cycle.js + 2012-08-19 + + Public Domain. + + NO WARRANTY EXPRESSED OR IMPLIED. USE AT YOUR OWN RISK. + + This code should be minified before deployment. + See http://javascript.crockford.com/jsmin.html + + USE YOUR OWN COPY. IT IS EXTREMELY UNWISE TO LOAD CODE FROM SERVERS YOU DO + NOT CONTROL. +*/ + +/*jslint evil: true, regexp: true */ + +/*members $ref, apply, call, decycle, hasOwnProperty, length, prototype, push, + retrocycle, stringify, test, toString +*/ + +if (typeof JSON.decycle !== 'function') { + JSON.decycle = function decycle(object) { + 'use strict'; + +// Make a deep copy of an object or array, assuring that there is at most +// one instance of each object or array in the resulting structure. The +// duplicate references (which might be forming cycles) are replaced with +// an object of the form +// {$ref: PATH} +// where the PATH is a JSONPath string that locates the first occurance. +// So, +// var a = []; +// a[0] = a; +// return JSON.stringify(JSON.decycle(a)); +// produces the string '[{"$ref":"$"}]'. + +// JSONPath is used to locate the unique object. $ indicates the top level of +// the object or array. [NUMBER] or [STRING] indicates a child member or +// property. + + var objects = [], // Keep a reference to each unique object or array + paths = []; // Keep the path to each unique object or array + + return (function derez(value, path) { + +// The derez recurses through the object, producing the deep copy. + + var i, // The loop counter + name, // Property name + nu; // The new object or array + + switch (typeof value) { + case 'object': + +// typeof null === 'object', so get out if this value is not really an object. +// Also get out if it is a weird builtin object. + + if (value === null || + value instanceof Boolean || + value instanceof Date || + value instanceof Number || + value instanceof RegExp || + value instanceof String) { + return value; + } + +// If the value is an object or array, look to see if we have already +// encountered it. If so, return a $ref/path object. This is a hard way, +// linear search that will get slower as the number of unique objects grows. + + for (i = 0; i < objects.length; i += 1) { + if (objects[i] === value) { + return {$ref: paths[i]}; + } + } + +// Otherwise, accumulate the unique value and its path. + + objects.push(value); + paths.push(path); + +// If it is an array, replicate the array. + + if (Object.prototype.toString.apply(value) === '[object Array]') { + nu = []; + for (i = 0; i < value.length; i += 1) { + nu[i] = derez(value[i], path + '[' + i + ']'); + } + } else { + +// If it is an object, replicate the object. + + nu = {}; + for (name in value) { + if (Object.prototype.hasOwnProperty.call(value, name)) { + nu[name] = derez(value[name], + path + '[' + JSON.stringify(name) + ']'); + } + } + } + return nu; + case 'number': + case 'string': + case 'boolean': + return value; + } + }(object, '$')); + }; +} + + +if (typeof JSON.retrocycle !== 'function') { + JSON.retrocycle = function retrocycle($) { + 'use strict'; + +// Restore an object that was reduced by decycle. Members whose values are +// objects of the form +// {$ref: PATH} +// are replaced with references to the value found by the PATH. This will +// restore cycles. The object will be mutated. + +// The eval function is used to locate the values described by a PATH. The +// root object is kept in a $ variable. A regular expression is used to +// assure that the PATH is extremely well formed. The regexp contains nested +// * quantifiers. That has been known to have extremely bad performance +// problems on some browsers for very long strings. A PATH is expected to be +// reasonably short. A PATH is allowed to belong to a very restricted subset of +// Goessner's JSONPath. + +// So, +// var s = '[{"$ref":"$"}]'; +// return JSON.retrocycle(JSON.parse(s)); +// produces an array containing a single element which is the array itself. + + var px = + /^\$(?:\[(?:\d+|\"(?:[^\\\"\u0000-\u001f]|\\([\\\"\/bfnrt]|u[0-9a-zA-Z]{4}))*\")\])*$/; + + (function rez(value) { + +// The rez function walks recursively through the object looking for $ref +// properties. When it finds one that has a value that is a path, then it +// replaces the $ref object with a reference to the value that is found by +// the path. + + var i, item, name, path; + + if (value && typeof value === 'object') { + if (Object.prototype.toString.apply(value) === '[object Array]') { + for (i = 0; i < value.length; i += 1) { + item = value[i]; + if (item && typeof item === 'object') { + path = item.$ref; + if (typeof path === 'string' && px.test(path)) { + value[i] = eval(path); + } else { + rez(item); + } + } + } + } else { + for (name in value) { + if (typeof value[name] === 'object') { + item = value[name]; + if (item) { + path = item.$ref; + if (typeof path === 'string' && px.test(path)) { + value[name] = eval(path); + } else { + rez(item); + } + } + } + } + } + } + }($)); + return $; + }; +} + +try { + module.exports = JSON; +} catch(e) { + +} \ No newline at end of file diff --git a/node_modules/json/index.js b/node_modules/json/index.js new file mode 100644 index 00000000..c7b632d4 --- /dev/null +++ b/node_modules/json/index.js @@ -0,0 +1,493 @@ +/* + json2.js + 2011-10-19 + + Public Domain. + + NO WARRANTY EXPRESSED OR IMPLIED. USE AT YOUR OWN RISK. + + See http://www.JSON.org/js.html + + + This code should be minified before deployment. + See http://javascript.crockford.com/jsmin.html + + USE YOUR OWN COPY. IT IS EXTREMELY UNWISE TO LOAD CODE FROM SERVERS YOU DO + NOT CONTROL. + + + This file creates a global JSON object containing two methods: stringify + and parse. + + JSON.stringify(value, replacer, space) + value any JavaScript value, usually an object or array. + + replacer an optional parameter that determines how object + values are stringified for objects. It can be a + function or an array of strings. + + space an optional parameter that specifies the indentation + of nested structures. If it is omitted, the text will + be packed without extra whitespace. If it is a number, + it will specify the number of spaces to indent at each + level. If it is a string (such as '\t' or ' '), + it contains the characters used to indent at each level. + + This method produces a JSON text from a JavaScript value. + + When an object value is found, if the object contains a toJSON + method, its toJSON method will be called and the result will be + stringified. A toJSON method does not serialize: it returns the + value represented by the name/value pair that should be serialized, + or undefined if nothing should be serialized. The toJSON method + will be passed the key associated with the value, and this will be + bound to the value + + For example, this would serialize Dates as ISO strings. + + Date.prototype.toJSON = function (key) { + function f(n) { + // Format integers to have at least two digits. + return n < 10 ? '0' + n : n; + } + + return this.getUTCFullYear() + '-' + + f(this.getUTCMonth() + 1) + '-' + + f(this.getUTCDate()) + 'T' + + f(this.getUTCHours()) + ':' + + f(this.getUTCMinutes()) + ':' + + f(this.getUTCSeconds()) + 'Z'; + }; + + You can provide an optional replacer method. It will be passed the + key and value of each member, with this bound to the containing + object. The value that is returned from your method will be + serialized. If your method returns undefined, then the member will + be excluded from the serialization. + + If the replacer parameter is an array of strings, then it will be + used to select the members to be serialized. It filters the results + such that only members with keys listed in the replacer array are + stringified. + + Values that do not have JSON representations, such as undefined or + functions, will not be serialized. Such values in objects will be + dropped; in arrays they will be replaced with null. You can use + a replacer function to replace those with JSON values. + JSON.stringify(undefined) returns undefined. + + The optional space parameter produces a stringification of the + value that is filled with line breaks and indentation to make it + easier to read. + + If the space parameter is a non-empty string, then that string will + be used for indentation. If the space parameter is a number, then + the indentation will be that many spaces. + + Example: + + text = JSON.stringify(['e', {pluribus: 'unum'}]); + // text is '["e",{"pluribus":"unum"}]' + + + text = JSON.stringify(['e', {pluribus: 'unum'}], null, '\t'); + // text is '[\n\t"e",\n\t{\n\t\t"pluribus": "unum"\n\t}\n]' + + text = JSON.stringify([new Date()], function (key, value) { + return this[key] instanceof Date ? + 'Date(' + this[key] + ')' : value; + }); + // text is '["Date(---current time---)"]' + + + JSON.parse(text, reviver) + This method parses a JSON text to produce an object or array. + It can throw a SyntaxError exception. + + The optional reviver parameter is a function that can filter and + transform the results. It receives each of the keys and values, + and its return value is used instead of the original value. + If it returns what it received, then the structure is not modified. + If it returns undefined then the member is deleted. + + Example: + + // Parse the text. Values that look like ISO date strings will + // be converted to Date objects. + + myData = JSON.parse(text, function (key, value) { + var a; + if (typeof value === 'string') { + a = +/^(\d{4})-(\d{2})-(\d{2})T(\d{2}):(\d{2}):(\d{2}(?:\.\d*)?)Z$/.exec(value); + if (a) { + return new Date(Date.UTC(+a[1], +a[2] - 1, +a[3], +a[4], + +a[5], +a[6])); + } + } + return value; + }); + + myData = JSON.parse('["Date(09/09/2001)"]', function (key, value) { + var d; + if (typeof value === 'string' && + value.slice(0, 5) === 'Date(' && + value.slice(-1) === ')') { + d = new Date(value.slice(5, -1)); + if (d) { + return d; + } + } + return value; + }); + + + This is a reference implementation. You are free to copy, modify, or + redistribute. +*/ + +/*jslint evil: true, regexp: true */ + +/*members "", "\b", "\t", "\n", "\f", "\r", "\"", JSON, "\\", apply, + call, charCodeAt, getUTCDate, getUTCFullYear, getUTCHours, + getUTCMinutes, getUTCMonth, getUTCSeconds, hasOwnProperty, join, + lastIndex, length, parse, prototype, push, replace, slice, stringify, + test, toJSON, toString, valueOf +*/ + + +// Create a JSON object only if one does not already exist. We create the +// methods in a closure to avoid creating global variables. + +var JSON; +if (!JSON) { + JSON = {}; +} + +(function () { + 'use strict'; + + function f(n) { + // Format integers to have at least two digits. + return n < 10 ? '0' + n : n; + } + + if (typeof Date.prototype.toJSON !== 'function') { + + Date.prototype.toJSON = function (key) { + + return isFinite(this.valueOf()) + ? this.getUTCFullYear() + '-' + + f(this.getUTCMonth() + 1) + '-' + + f(this.getUTCDate()) + 'T' + + f(this.getUTCHours()) + ':' + + f(this.getUTCMinutes()) + ':' + + f(this.getUTCSeconds()) + 'Z' + : null; + }; + + String.prototype.toJSON = + Number.prototype.toJSON = + Boolean.prototype.toJSON = function (key) { + return this.valueOf(); + }; + } + + var cx = /[\u0000\u00ad\u0600-\u0604\u070f\u17b4\u17b5\u200c-\u200f\u2028-\u202f\u2060-\u206f\ufeff\ufff0-\uffff]/g, + escapable = /[\\\"\x00-\x1f\x7f-\x9f\u00ad\u0600-\u0604\u070f\u17b4\u17b5\u200c-\u200f\u2028-\u202f\u2060-\u206f\ufeff\ufff0-\uffff]/g, + gap, + indent, + meta = { // table of character substitutions + '\b': '\\b', + '\t': '\\t', + '\n': '\\n', + '\f': '\\f', + '\r': '\\r', + '"' : '\\"', + '\\': '\\\\' + }, + rep; + + + function quote(string) { + +// If the string contains no control characters, no quote characters, and no +// backslash characters, then we can safely slap some quotes around it. +// Otherwise we must also replace the offending characters with safe escape +// sequences. + + escapable.lastIndex = 0; + return escapable.test(string) ? '"' + string.replace(escapable, function (a) { + var c = meta[a]; + return typeof c === 'string' + ? c + : '\\u' + ('0000' + a.charCodeAt(0).toString(16)).slice(-4); + }) + '"' : '"' + string + '"'; + } + + + function str(key, holder) { + +// Produce a string from holder[key]. + + var i, // The loop counter. + k, // The member key. + v, // The member value. + length, + mind = gap, + partial, + value = holder[key]; + +// If the value has a toJSON method, call it to obtain a replacement value. + + if (value && typeof value === 'object' && + typeof value.toJSON === 'function') { + value = value.toJSON(key); + } + +// If we were called with a replacer function, then call the replacer to +// obtain a replacement value. + + if (typeof rep === 'function') { + value = rep.call(holder, key, value); + } + +// What happens next depends on the value's type. + + switch (typeof value) { + case 'string': + return quote(value); + + case 'number': + +// JSON numbers must be finite. Encode non-finite numbers as null. + + return isFinite(value) ? String(value) : 'null'; + + case 'boolean': + case 'null': + +// If the value is a boolean or null, convert it to a string. Note: +// typeof null does not produce 'null'. The case is included here in +// the remote chance that this gets fixed someday. + + return String(value); + +// If the type is 'object', we might be dealing with an object or an array or +// null. + + case 'object': + +// Due to a specification blunder in ECMAScript, typeof null is 'object', +// so watch out for that case. + + if (!value) { + return 'null'; + } + +// Make an array to hold the partial results of stringifying this object value. + + gap += indent; + partial = []; + +// Is the value an array? + + if (Object.prototype.toString.apply(value) === '[object Array]') { + +// The value is an array. Stringify every element. Use null as a placeholder +// for non-JSON values. + + length = value.length; + for (i = 0; i < length; i += 1) { + partial[i] = str(i, value) || 'null'; + } + +// Join all of the elements together, separated with commas, and wrap them in +// brackets. + + v = partial.length === 0 + ? '[]' + : gap + ? '[\n' + gap + partial.join(',\n' + gap) + '\n' + mind + ']' + : '[' + partial.join(',') + ']'; + gap = mind; + return v; + } + +// If the replacer is an array, use it to select the members to be stringified. + + if (rep && typeof rep === 'object') { + length = rep.length; + for (i = 0; i < length; i += 1) { + if (typeof rep[i] === 'string') { + k = rep[i]; + v = str(k, value); + if (v) { + partial.push(quote(k) + (gap ? ': ' : ':') + v); + } + } + } + } else { + +// Otherwise, iterate through all of the keys in the object. + + for (k in value) { + if (Object.prototype.hasOwnProperty.call(value, k)) { + v = str(k, value); + if (v) { + partial.push(quote(k) + (gap ? ': ' : ':') + v); + } + } + } + } + +// Join all of the member texts together, separated with commas, +// and wrap them in braces. + + v = partial.length === 0 + ? '{}' + : gap + ? '{\n' + gap + partial.join(',\n' + gap) + '\n' + mind + '}' + : '{' + partial.join(',') + '}'; + gap = mind; + return v; + } + } + +// If the JSON object does not yet have a stringify method, give it one. + + if (typeof JSON.stringify !== 'function') { + JSON.stringify = function (value, replacer, space) { + +// The stringify method takes a value and an optional replacer, and an optional +// space parameter, and returns a JSON text. The replacer can be a function +// that can replace values, or an array of strings that will select the keys. +// A default replacer method can be provided. Use of the space parameter can +// produce text that is more easily readable. + + var i; + gap = ''; + indent = ''; + +// If the space parameter is a number, make an indent string containing that +// many spaces. + + if (typeof space === 'number') { + for (i = 0; i < space; i += 1) { + indent += ' '; + } + +// If the space parameter is a string, it will be used as the indent string. + + } else if (typeof space === 'string') { + indent = space; + } + +// If there is a replacer, it must be a function or an array. +// Otherwise, throw an error. + + rep = replacer; + if (replacer && typeof replacer !== 'function' && + (typeof replacer !== 'object' || + typeof replacer.length !== 'number')) { + throw new Error('JSON.stringify'); + } + +// Make a fake root object containing our value under the key of ''. +// Return the result of stringifying the value. + + return str('', {'': value}); + }; + } + + +// If the JSON object does not yet have a parse method, give it one. + + if (typeof JSON.parse !== 'function') { + JSON.parse = function (text, reviver) { + +// The parse method takes a text and an optional reviver function, and returns +// a JavaScript value if the text is a valid JSON text. + + var j; + + function walk(holder, key) { + +// The walk method is used to recursively walk the resulting structure so +// that modifications can be made. + + var k, v, value = holder[key]; + if (value && typeof value === 'object') { + for (k in value) { + if (Object.prototype.hasOwnProperty.call(value, k)) { + v = walk(value, k); + if (v !== undefined) { + value[k] = v; + } else { + delete value[k]; + } + } + } + } + return reviver.call(holder, key, value); + } + + +// Parsing happens in four stages. In the first stage, we replace certain +// Unicode characters with escape sequences. JavaScript handles many characters +// incorrectly, either silently deleting them, or treating them as line endings. + + text = String(text); + cx.lastIndex = 0; + if (cx.test(text)) { + text = text.replace(cx, function (a) { + return '\\u' + + ('0000' + a.charCodeAt(0).toString(16)).slice(-4); + }); + } + +// In the second stage, we run the text against regular expressions that look +// for non-JSON patterns. We are especially concerned with '()' and 'new' +// because they can cause invocation, and '=' because it can cause mutation. +// But just to be safe, we want to reject all unexpected forms. + +// We split the second stage into 4 regexp operations in order to work around +// crippling inefficiencies in IE's and Safari's regexp engines. First we +// replace the JSON backslash pairs with '@' (a non-JSON character). Second, we +// replace all simple value tokens with ']' characters. Third, we delete all +// open brackets that follow a colon or comma or that begin the text. Finally, +// we look to see that the remaining characters are only whitespace or ']' or +// ',' or ':' or '{' or '}'. If that is so, then the text is safe for eval. + + if (/^[\],:{}\s]*$/ + .test(text.replace(/\\(?:["\\\/bfnrt]|u[0-9a-fA-F]{4})/g, '@') + .replace(/"[^"\\\n\r]*"|true|false|null|-?\d+(?:\.\d*)?(?:[eE][+\-]?\d+)?/g, ']') + .replace(/(?:^|:|,)(?:\s*\[)+/g, ''))) { + +// In the third stage we use the eval function to compile the text into a +// JavaScript structure. The '{' operator is subject to a syntactic ambiguity +// in JavaScript: it can begin a block or an object literal. We wrap the text +// in parens to eliminate the ambiguity. + + j = eval('(' + text + ')'); + +// In the optional fourth stage, we recursively walk the new structure, passing +// each name/value pair to a reviver function for possible transformation. + + return typeof reviver === 'function' + ? walk({'': j}, '') + : j; + } + +// If the text is not JSON parseable, then a SyntaxError is thrown. + + throw new SyntaxError('JSON.parse'); + }; + } +}()); + +try { + module.exports = JSON; +} catch(e) { + +} diff --git a/package.json b/package.json new file mode 100644 index 00000000..8a8cccd9 --- /dev/null +++ b/package.json @@ -0,0 +1,30 @@ +{ + "name": "socketcluster", + "description": "SocketCluster - A Highly parallelized WebSocket server cluster to make the most of multi-core machines/instances.", + "version": "0.9.0", + "homepage": "https://github.com/topcloud/socketcluster", + "contributors": [ + { + "name": "Jonathan Gros-Dubois", + "email": "grosjona@yahoo.com.au" + } + ], + "dependencies": { + "iocluster": ">=0.9.13", + "loadbalancer": ">=0.9.9", + "ndata": ">=0.9.33", + "socketcluster-server": ">=0.9.12" + }, + "bundleDependencies": [ + "json" + ], + "keywords": [ + "websocket", + "server", + "engine.io", + "cluster", + "distributed", + "parallelized" + ], + "readmeFilename": "README.md" +} diff --git a/socketworker.js b/socketworker.js new file mode 100644 index 00000000..54086527 --- /dev/null +++ b/socketworker.js @@ -0,0 +1,205 @@ +var socketClusterServer = require('socketcluster-server'); +var EventEmitter = require('events').EventEmitter; +var crypto = require('crypto'); +var domain = require('domain'); +var http = require('http'); + +var SocketWorker = function (options) { + var self = this; + + this._errorDomain = domain.create(); + this._errorDomain.on('error', function () { + self.errorHandler.apply(self, arguments); + if (self._options.rebootWorkerOnError) { + self.emit('exit'); + } + }); + + this.start = this._errorDomain.bind(this._start); + this._errorDomain.run(function () { + self._init(options); + }); +}; + +SocketWorker.prototype = Object.create(EventEmitter.prototype); + +SocketWorker.prototype._init = function (options) { + var self = this; + + this._options = { + transports: ['polling', 'websocket'], + host: 'localhost' + }; + + for (var i in options) { + this._options[i] = options[i]; + } + + if (this._options.dataKey == null) { + this._options.dataKey = crypto.randomBytes(32).toString('hex'); + } + + this._clusterEngine = require(this._options.clusterEngine); + + this._paths = options.paths; + + this._httpRequestCount = 0; + this._ioRequestCount = 0; + this._httpRPM = 0; + this._ioRPM = 0; + + this._ioClusterClient = new this._clusterEngine.IOClusterClient({ + stores: this._options.stores, + dataKey: this._options.dataKey, + connectTimeout: this._options.connectTimeout, + dataExpiry: this._options.sessionTimeout, + heartRate: this._options.sessionHeartRate, + addressSocketLimit: this._options.addressSocketLimit + }); + + this._errorDomain.add(this._ioClusterClient); + + this._server = http.createServer(); + this._errorDomain.add(this._server); + + this._socketServer = socketClusterServer.attach(this._server, { + sourcePort: this._options.sourcePort, + ioClusterClient: this._ioClusterClient, + transports: this._options.transports, + pingTimeout: this._options.heartbeatTimeout, + pingInterval: this._options.heartbeatInterval, + upgradeTimeout: this._options.connectTimeout, + host: this._options.host, + secure: this._options.protocol == 'https', + appName: this._options.appName + }); + + this._socketServer.on('connection', function (socket) { + socket.on('message', function () { + self._ioRequestCount++; + }); + self.emit('connection', socket); + }); + + this._socketURL = this._socketServer.getURL(); + this._socketURLRegex = new RegExp('^' + this._socketURL); + + this._errorDomain.add(this._socketServer); + this._socketServer.on('ready', function () { + self._server.on('request', self._httpRequestHandler.bind(self)); + self.emit('ready'); + }); +}; + +SocketWorker.prototype.getSocketURL = function () { + return this._socketURL; +}; + +SocketWorker.prototype._start = function () { + this._httpRequestCount = 0; + this._ioRequestCount = 0; + this._httpRPM = 0; + this._ioRPM = 0; + + if (this._statusInterval != null) { + clearInterval(this._statusInterval); + } + this._statusInterval = setInterval(this._calculateStatus.bind(this), this._options.workerStatusInterval * 1000); + + this._server.listen(this._options.workerPort); +}; + +SocketWorker.prototype._httpRequestHandler = function (req, res) { + this._httpRequestCount++; + if (req.url == this._paths.statusURL) { + this._handleStatusRequest(req, res); + } else if (!this._socketURLRegex.test(req.url)) { + this._server.emit('req', req, res); + } +}; + +SocketWorker.prototype._handleStatusRequest = function (req, res) { + var self = this; + + var isOpen = true; + + var reqTimeout = setTimeout(function () { + if (isOpen) { + res.writeHead(500, { + 'Content-Type': 'application/json' + }); + res.end(); + isOpen = false; + } + }, this._options.connectTimeout * 1000); + + var buffers = []; + req.on('data', function (chunk) { + buffers.push(chunk); + }); + + req.on('end', function () { + clearTimeout(reqTimeout); + if (isOpen) { + var statusReq = null; + try { + statusReq = JSON.parse(Buffer.concat(buffers).toString()); + } catch (e) {} + + if (statusReq && statusReq.dataKey == self._options.dataKey) { + var status = JSON.stringify(self.getStatus()); + res.writeHead(200, { + 'Content-Type': 'application/json' + }); + res.end(status); + } else { + res.writeHead(401, { + 'Content-Type': 'application/json' + }); + res.end(); + } + isOpen = false; + } + }); +}; + +SocketWorker.prototype.getWSServer = function () { + return this._socketServer; +}; + +SocketWorker.prototype.getHTTPServer = function () { + return this._server; +}; + +SocketWorker.prototype._calculateStatus = function () { + var perMinuteFactor = 60 / this._options.workerStatusInterval; + this._httpRPM = this._httpRequestCount * perMinuteFactor; + this._ioRPM = this._ioRequestCount * perMinuteFactor; + this._httpRequestCount = 0; + this._ioRequestCount = 0; +}; + +SocketWorker.prototype.getStatus = function () { + return { + clientCount: this._socketServer.clientsCount, + httpRPM: this._httpRPM, + ioRPM: this._ioRPM + }; +}; + +SocketWorker.prototype.handleMasterEvent = function () { + this.emit.apply(this, arguments); +}; + +SocketWorker.prototype.errorHandler = function (err) { + this.emit('error', err); +}; + +SocketWorker.prototype.noticeHandler = function (notice) { + if (notice.message != null) { + notice = notice.message; + } + this.emit('notice', notice); +}; + +module.exports = SocketWorker; \ No newline at end of file diff --git a/worker-bootstrap.js b/worker-bootstrap.js new file mode 100644 index 00000000..0c9d34fb --- /dev/null +++ b/worker-bootstrap.js @@ -0,0 +1,55 @@ +var SocketWorker = require('./socketworker'); +var worker; + +var handleError = function (err) { + var error; + if (err.stack) { + error = { + message: err.message, + stack: err.stack + } + } else { + error = err; + } + process.send({type: 'error', data: error}); +}; + +var handleNotice = function (notice) { + if (notice instanceof Error) { + notice = notice.message; + } + process.send({type: 'notice', data: notice}); +}; + +var handleReady = function () { + process.send({type: 'ready'}); +}; + +var handleExit = function () { + process.exit(); +}; + +process.on('message', function (m) { + if (m.type == 'init') { + worker = new SocketWorker(m.data); + + if (m.data.propagateErrors) { + worker.on('error', handleError); + worker.on('notice', handleNotice); + worker.on('ready', handleReady); + worker.on('exit', handleExit); + } + + var workerController = require(m.data.paths.appWorkerControllerPath); + worker.on('ready', function () { + workerController.run(worker); + worker.start(); + }); + } else if (m.type == 'emit') { + if (m.data) { + worker.handleMasterEvent(m.event, m.data); + } else { + worker.handleMasterEvent(m.event); + } + } +}); \ No newline at end of file From 4d27143a0a2dd22dd391d6bfab20859dc04b9d02 Mon Sep 17 00:00:00 2001 From: Jonathan Gros-Dubois Date: Sun, 6 Apr 2014 21:03:33 +1000 Subject: [PATCH 0003/1219] Fixed readme --- README.md | 51 ++++++++++++++++++++++++++------------------------- 1 file changed, 26 insertions(+), 25 deletions(-) diff --git a/README.md b/README.md index 9120119b..ff37828f 100644 --- a/README.md +++ b/README.md @@ -136,33 +136,34 @@ module.exports.run = function (worker) { Exposed by `require('socketcluster').SocketCluster`. ### SocketCluster(opts:Object) - Creates a new SocketCluster, must be invoked with the new keyword. - - ```js - var SocketCluster = require('socketcluster').SocketCluster; - - var socketCluster = new SocketCluster({ - workers: [9100, 9101, 9102], - stores: [9001, 9002, 9003], - port: 8000, - appName: 'myapp', - workerController: 'worker.js' - }); - ``` - - Documentation on all supported options is coming soon (there are around 30 of them - Most of them are optional). + +Creates a new SocketCluster, must be invoked with the new keyword. + +```js +var SocketCluster = require('socketcluster').SocketCluster; + +var socketCluster = new SocketCluster({ + workers: [9100, 9101, 9102], + stores: [9001, 9002, 9003], + port: 8000, + appName: 'myapp', + workerController: 'worker.js' +}); +``` + +Documentation on all supported options is coming soon (there are around 30 of them - Most of them are optional). ### SocketWorker - - A SocketWorker object is passed as the argument to your workerController's run(worker) function. - Example - Inside worker.js: - - ```js - module.exports.run = function (worker) { - // worker here is an instance of SocketWorker - } - ``` + +A SocketWorker object is passed as the argument to your workerController's run(worker) function. +Example - Inside worker.js: + +```js +module.exports.run = function (worker) { + // worker here is an instance of SocketWorker +} +``` ### ClusterServer - A ClusterServer instance is returned from worker.getWSServer() - You use it to handle WebSocket connections. \ No newline at end of file +A ClusterServer instance is returned from worker.getWSServer() - You use it to handle WebSocket connections. \ No newline at end of file From 932a3819037824debbb2ccd36fe1277e3a12120d Mon Sep 17 00:00:00 2001 From: Jonathan Gros-Dubois Date: Sun, 6 Apr 2014 21:05:15 +1000 Subject: [PATCH 0004/1219] Fixed readme some more --- README.md | 34 +++++++++++++++++----------------- 1 file changed, 17 insertions(+), 17 deletions(-) diff --git a/README.md b/README.md index ff37828f..ab4a4d1c 100644 --- a/README.md +++ b/README.md @@ -115,25 +115,25 @@ Using SocketCluster with express is simple, you put the code inside your workerC ```js module.exports.run = function (worker) { - // Get a reference to our raw Node HTTP server + // Get a reference to our raw Node HTTP server var httpServer = worker.getHTTPServer(); // Get a reference to our WebSocket server var wsServer = worker.getWSServer(); - - var app = require('express')(); - - // Add whatever express middleware you like... - - // Make your express app handle all essential requests - httpServer.on('req', app); + + var app = require('express')(); + + // Add whatever express middleware you like... + + // Make your express app handle all essential requests + httpServer.on('req', app); } ``` -## API (Coming soon) +## API (Documentation coming soon) ### SocketCluster - Exposed by `require('socketcluster').SocketCluster`. +Exposed by `require('socketcluster').SocketCluster`. ### SocketCluster(opts:Object) @@ -143,16 +143,16 @@ Creates a new SocketCluster, must be invoked with the new keyword. var SocketCluster = require('socketcluster').SocketCluster; var socketCluster = new SocketCluster({ - workers: [9100, 9101, 9102], - stores: [9001, 9002, 9003], - port: 8000, - appName: 'myapp', - workerController: 'worker.js' + workers: [9100, 9101, 9102], + stores: [9001, 9002, 9003], + port: 8000, + appName: 'myapp', + workerController: 'worker.js' }); ``` Documentation on all supported options is coming soon (there are around 30 of them - Most of them are optional). - + ### SocketWorker A SocketWorker object is passed as the argument to your workerController's run(worker) function. @@ -160,7 +160,7 @@ Example - Inside worker.js: ```js module.exports.run = function (worker) { - // worker here is an instance of SocketWorker + // worker here is an instance of SocketWorker } ``` From 2c14f363a5a47889128b08ec22d98bc4cb455027 Mon Sep 17 00:00:00 2001 From: Jonathan Gros-Dubois Date: Sun, 6 Apr 2014 21:19:02 +1000 Subject: [PATCH 0005/1219] Additional info to setup socketcluster-client --- README.md | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/README.md b/README.md index ab4a4d1c..985f457c 100644 --- a/README.md +++ b/README.md @@ -22,6 +22,15 @@ To install, run: npm install socketcluster ``` +Note that to use socketcluster you will also need the client which you an get using the following command: + +```bash +npm install socketcluster-client +``` + +The socketcluster-client script is called socketcluster.js (located in the main socketcluster-client directory) +- You should include it in your HTML page using a + + + + + \ No newline at end of file diff --git a/sample/public/socketcluster.js b/sample/public/socketcluster.js new file mode 100644 index 00000000..fecabe81 --- /dev/null +++ b/sample/public/socketcluster.js @@ -0,0 +1,4369 @@ +!function(e){if("object"==typeof exports&&"undefined"!=typeof module)module.exports=e();else if("function"==typeof define&&define.amd)define([],e);else{var f;"undefined"!=typeof window?f=window:"undefined"!=typeof global?f=global:"undefined"!=typeof self&&(f=self),f.socketCluster=e()}}(function(){var define,module,exports;return (function e(t,n,r){function s(o,u){if(!n[o]){if(!t[o]){var a=typeof require=="function"&&require;if(!u&&a)return a(o,!0);if(i)return i(o,!0);throw new Error("Cannot find module '"+o+"'")}var f=n[o]={exports:{}};t[o][0].call(f.exports,function(e){var n=t[o][1][e];return s(n?n:e)},f,f.exports,e,t,n,r)}return n[o].exports}var i=typeof require=="function"&&require;for(var o=0;o + * MIT Licensed + */ + +/** + * Based on JSON2 (http://www.JSON.org/js.html). + */ + +(function (exports, nativeJSON) { + "use strict"; + + // use native JSON if it's available + if (nativeJSON && nativeJSON.parse){ + return exports.JSON = { + parse: nativeJSON.parse + , stringify: nativeJSON.stringify + }; + } + + var JSON = exports.JSON = {}; + + function f(n) { + // Format integers to have at least two digits. + return n < 10 ? '0' + n : n; + } + + function date(d, key) { + return isFinite(d.valueOf()) ? + d.getUTCFullYear() + '-' + + f(d.getUTCMonth() + 1) + '-' + + f(d.getUTCDate()) + 'T' + + f(d.getUTCHours()) + ':' + + f(d.getUTCMinutes()) + ':' + + f(d.getUTCSeconds()) + 'Z' : null; + }; + + var cx = /[\u0000\u00ad\u0600-\u0604\u070f\u17b4\u17b5\u200c-\u200f\u2028-\u202f\u2060-\u206f\ufeff\ufff0-\uffff]/g, + escapable = /[\\\"\x00-\x1f\x7f-\x9f\u00ad\u0600-\u0604\u070f\u17b4\u17b5\u200c-\u200f\u2028-\u202f\u2060-\u206f\ufeff\ufff0-\uffff]/g, + gap, + indent, + meta = { // table of character substitutions + '\b': '\\b', + '\t': '\\t', + '\n': '\\n', + '\f': '\\f', + '\r': '\\r', + '"' : '\\"', + '\\': '\\\\' + }, + rep; + + + function quote(string) { + +// If the string contains no control characters, no quote characters, and no +// backslash characters, then we can safely slap some quotes around it. +// Otherwise we must also replace the offending characters with safe escape +// sequences. + + escapable.lastIndex = 0; + return escapable.test(string) ? '"' + string.replace(escapable, function (a) { + var c = meta[a]; + return typeof c === 'string' ? c : + '\\u' + ('0000' + a.charCodeAt(0).toString(16)).slice(-4); + }) + '"' : '"' + string + '"'; + } + + + function str(key, holder) { + +// Produce a string from holder[key]. + + var i, // The loop counter. + k, // The member key. + v, // The member value. + length, + mind = gap, + partial, + value = holder[key]; + +// If the value has a toJSON method, call it to obtain a replacement value. + + if (value instanceof Date) { + value = date(key); + } + +// If we were called with a replacer function, then call the replacer to +// obtain a replacement value. + + if (typeof rep === 'function') { + value = rep.call(holder, key, value); + } + +// What happens next depends on the value's type. + + switch (typeof value) { + case 'string': + return quote(value); + + case 'number': + +// JSON numbers must be finite. Encode non-finite numbers as null. + + return isFinite(value) ? String(value) : 'null'; + + case 'boolean': + case 'null': + +// If the value is a boolean or null, convert it to a string. Note: +// typeof null does not produce 'null'. The case is included here in +// the remote chance that this gets fixed someday. + + return String(value); + +// If the type is 'object', we might be dealing with an object or an array or +// null. + + case 'object': + +// Due to a specification blunder in ECMAScript, typeof null is 'object', +// so watch out for that case. + + if (!value) { + return 'null'; + } + +// Make an array to hold the partial results of stringifying this object value. + + gap += indent; + partial = []; + +// Is the value an array? + + if (Object.prototype.toString.apply(value) === '[object Array]') { + +// The value is an array. Stringify every element. Use null as a placeholder +// for non-JSON values. + + length = value.length; + for (i = 0; i < length; i += 1) { + partial[i] = str(i, value) || 'null'; + } + +// Join all of the elements together, separated with commas, and wrap them in +// brackets. + + v = partial.length === 0 ? '[]' : gap ? + '[\n' + gap + partial.join(',\n' + gap) + '\n' + mind + ']' : + '[' + partial.join(',') + ']'; + gap = mind; + return v; + } + +// If the replacer is an array, use it to select the members to be stringified. + + if (rep && typeof rep === 'object') { + length = rep.length; + for (i = 0; i < length; i += 1) { + if (typeof rep[i] === 'string') { + k = rep[i]; + v = str(k, value); + if (v) { + partial.push(quote(k) + (gap ? ': ' : ':') + v); + } + } + } + } else { + +// Otherwise, iterate through all of the keys in the object. + + for (k in value) { + if (Object.prototype.hasOwnProperty.call(value, k)) { + v = str(k, value); + if (v) { + partial.push(quote(k) + (gap ? ': ' : ':') + v); + } + } + } + } + +// Join all of the member texts together, separated with commas, +// and wrap them in braces. + + v = partial.length === 0 ? '{}' : gap ? + '{\n' + gap + partial.join(',\n' + gap) + '\n' + mind + '}' : + '{' + partial.join(',') + '}'; + gap = mind; + return v; + } + } + +// If the JSON object does not yet have a stringify method, give it one. + + JSON.stringify = function (value, replacer, space) { + +// The stringify method takes a value and an optional replacer, and an optional +// space parameter, and returns a JSON text. The replacer can be a function +// that can replace values, or an array of strings that will select the keys. +// A default replacer method can be provided. Use of the space parameter can +// produce text that is more easily readable. + + var i; + gap = ''; + indent = ''; + +// If the space parameter is a number, make an indent string containing that +// many spaces. + + if (typeof space === 'number') { + for (i = 0; i < space; i += 1) { + indent += ' '; + } + +// If the space parameter is a string, it will be used as the indent string. + + } else if (typeof space === 'string') { + indent = space; + } + +// If there is a replacer, it must be a function or an array. +// Otherwise, throw an error. + + rep = replacer; + if (replacer && typeof replacer !== 'function' && + (typeof replacer !== 'object' || + typeof replacer.length !== 'number')) { + throw new Error('JSON.stringify'); + } + +// Make a fake root object containing our value under the key of ''. +// Return the result of stringifying the value. + + return str('', {'': value}); + }; + +// If the JSON object does not yet have a parse method, give it one. + + JSON.parse = function (text, reviver) { + // The parse method takes a text and an optional reviver function, and returns + // a JavaScript value if the text is a valid JSON text. + + var j; + + function walk(holder, key) { + + // The walk method is used to recursively walk the resulting structure so + // that modifications can be made. + + var k, v, value = holder[key]; + if (value && typeof value === 'object') { + for (k in value) { + if (Object.prototype.hasOwnProperty.call(value, k)) { + v = walk(value, k); + if (v !== undefined) { + value[k] = v; + } else { + delete value[k]; + } + } + } + } + return reviver.call(holder, key, value); + } + + + // Parsing happens in four stages. In the first stage, we replace certain + // Unicode characters with escape sequences. JavaScript handles many characters + // incorrectly, either silently deleting them, or treating them as line endings. + + text = String(text); + cx.lastIndex = 0; + if (cx.test(text)) { + text = text.replace(cx, function (a) { + return '\\u' + + ('0000' + a.charCodeAt(0).toString(16)).slice(-4); + }); + } + + // In the second stage, we run the text against regular expressions that look + // for non-JSON patterns. We are especially concerned with '()' and 'new' + // because they can cause invocation, and '=' because it can cause mutation. + // But just to be safe, we want to reject all unexpected forms. + + // We split the second stage into 4 regexp operations in order to work around + // crippling inefficiencies in IE's and Safari's regexp engines. First we + // replace the JSON backslash pairs with '@' (a non-JSON character). Second, we + // replace all simple value tokens with ']' characters. Third, we delete all + // open brackets that follow a colon or comma or that begin the text. Finally, + // we look to see that the remaining characters are only whitespace or ']' or + // ',' or ':' or '{' or '}'. If that is so, then the text is safe for eval. + + if (/^[\],:{}\s]*$/ + .test(text.replace(/\\(?:["\\\/bfnrt]|u[0-9a-fA-F]{4})/g, '@') + .replace(/"[^"\\\n\r]*"|true|false|null|-?\d+(?:\.\d*)?(?:[eE][+\-]?\d+)?/g, ']') + .replace(/(?:^|:|,)(?:\s*\[)+/g, ''))) { + + // In the third stage we use the eval function to compile the text into a + // JavaScript structure. The '{' operator is subject to a syntactic ambiguity + // in JavaScript: it can begin a block or an object literal. We wrap the text + // in parens to eliminate the ambiguity. + + j = eval('(' + text + ')'); + + // In the optional fourth stage, we recursively walk the new structure, passing + // each name/value pair to a reviver function for possible transformation. + + return typeof reviver === 'function' ? + walk({'': j}, '') : j; + } + + // If the text is not JSON parseable, then a SyntaxError is thrown. + + throw new SyntaxError('JSON.parse'); + }; + +})( + 'undefined' != typeof io ? io : module.exports + , typeof JSON !== 'undefined' ? JSON : undefined +); +},{}],3:[function(_dereq_,module,exports){ + +/** + * Module dependencies. + */ + +var index = _dereq_('indexof'); + +/** + * Expose `Emitter`. + */ + +module.exports = Emitter; + +/** + * Initialize a new `Emitter`. + * + * @api public + */ + +function Emitter(obj) { + if (obj) return mixin(obj); +}; + +/** + * Mixin the emitter properties. + * + * @param {Object} obj + * @return {Object} + * @api private + */ + +function mixin(obj) { + for (var key in Emitter.prototype) { + obj[key] = Emitter.prototype[key]; + } + return obj; +} + +/** + * Listen on the given `event` with `fn`. + * + * @param {String} event + * @param {Function} fn + * @return {Emitter} + * @api public + */ + +Emitter.prototype.on = function(event, fn){ + this._callbacks = this._callbacks || {}; + (this._callbacks[event] = this._callbacks[event] || []) + .push(fn); + return this; +}; + +/** + * Adds an `event` listener that will be invoked a single + * time then automatically removed. + * + * @param {String} event + * @param {Function} fn + * @return {Emitter} + * @api public + */ + +Emitter.prototype.once = function(event, fn){ + var self = this; + this._callbacks = this._callbacks || {}; + + function on() { + self.off(event, on); + fn.apply(this, arguments); + } + + fn._off = on; + this.on(event, on); + return this; +}; + +/** + * Remove the given callback for `event` or all + * registered callbacks. + * + * @param {String} event + * @param {Function} fn + * @return {Emitter} + * @api public + */ + +Emitter.prototype.off = +Emitter.prototype.removeListener = +Emitter.prototype.removeAllListeners = function(event, fn){ + this._callbacks = this._callbacks || {}; + + // all + if (0 == arguments.length) { + this._callbacks = {}; + return this; + } + + // specific event + var callbacks = this._callbacks[event]; + if (!callbacks) return this; + + // remove all handlers + if (1 == arguments.length) { + delete this._callbacks[event]; + return this; + } + + // remove specific handler + var i = index(callbacks, fn._off || fn); + if (~i) callbacks.splice(i, 1); + return this; +}; + +/** + * Emit `event` with the given args. + * + * @param {String} event + * @param {Mixed} ... + * @return {Emitter} + */ + +Emitter.prototype.emit = function(event){ + this._callbacks = this._callbacks || {}; + var args = [].slice.call(arguments, 1) + , callbacks = this._callbacks[event]; + + if (callbacks) { + callbacks = callbacks.slice(0); + for (var i = 0, len = callbacks.length; i < len; ++i) { + callbacks[i].apply(this, args); + } + } + + return this; +}; + +/** + * Return array of callbacks for `event`. + * + * @param {String} event + * @return {Array} + * @api public + */ + +Emitter.prototype.listeners = function(event){ + this._callbacks = this._callbacks || {}; + return this._callbacks[event] || []; +}; + +/** + * Check if this emitter has `event` handlers. + * + * @param {String} event + * @return {Boolean} + * @api public + */ + +Emitter.prototype.hasListeners = function(event){ + return !! this.listeners(event).length; +}; + +},{"indexof":4}],4:[function(_dereq_,module,exports){ + +var indexOf = [].indexOf; + +module.exports = function(arr, obj){ + if (indexOf) return arr.indexOf(obj); + for (var i = 0; i < arr.length; ++i) { + if (arr[i] === obj) return i; + } + return -1; +}; +},{}],5:[function(_dereq_,module,exports){ + +module.exports = _dereq_('./lib/'); + +},{"./lib/":6}],6:[function(_dereq_,module,exports){ + +module.exports = _dereq_('./socket'); + +/** + * Exports parser + * + * @api public + * + */ +module.exports.parser = _dereq_('engine.io-parser'); + +},{"./socket":7,"engine.io-parser":19}],7:[function(_dereq_,module,exports){ +(function (global){ +/** + * Module dependencies. + */ + +var transports = _dereq_('./transports'); +var Emitter = _dereq_('component-emitter'); +var debug = _dereq_('debug')('engine.io-client:socket'); +var index = _dereq_('indexof'); +var parser = _dereq_('engine.io-parser'); +var parseuri = _dereq_('parseuri'); +var parsejson = _dereq_('parsejson'); +var parseqs = _dereq_('parseqs'); + +/** + * Module exports. + */ + +module.exports = Socket; + +/** + * Noop function. + * + * @api private + */ + +function noop(){} + +/** + * Socket constructor. + * + * @param {String|Object} uri or options + * @param {Object} options + * @api public + */ + +function Socket(uri, opts){ + if (!(this instanceof Socket)) return new Socket(uri, opts); + + opts = opts || {}; + + if (uri && 'object' == typeof uri) { + opts = uri; + uri = null; + } + + if (uri) { + uri = parseuri(uri); + opts.host = uri.host; + opts.secure = uri.protocol == 'https' || uri.protocol == 'wss'; + opts.port = uri.port; + if (uri.query) opts.query = uri.query; + } + + this.secure = null != opts.secure ? opts.secure : + (global.location && 'https:' == location.protocol); + + if (opts.host) { + var pieces = opts.host.split(':'); + opts.hostname = pieces.shift(); + if (pieces.length) opts.port = pieces.pop(); + } + + this.agent = opts.agent || false; + this.hostname = opts.hostname || + (global.location ? location.hostname : 'localhost'); + this.port = opts.port || (global.location && location.port ? + location.port : + (this.secure ? 443 : 80)); + this.query = opts.query || {}; + if ('string' == typeof this.query) this.query = parseqs.decode(this.query); + this.upgrade = false !== opts.upgrade; + this.path = (opts.path || '/engine.io').replace(/\/$/, '') + '/'; + this.forceJSONP = !!opts.forceJSONP; + this.forceBase64 = !!opts.forceBase64; + this.timestampParam = opts.timestampParam || 't'; + this.timestampRequests = opts.timestampRequests; + this.transports = opts.transports || ['polling', 'websocket']; + this.readyState = ''; + this.writeBuffer = []; + this.callbackBuffer = []; + this.policyPort = opts.policyPort || 843; + this.rememberUpgrade = opts.rememberUpgrade || false; + this.open(); + this.binaryType = null; + this.onlyBinaryUpgrades = opts.onlyBinaryUpgrades; +} + +Socket.priorWebsocketSuccess = false; + +/** + * Mix in `Emitter`. + */ + +Emitter(Socket.prototype); + +/** + * Protocol version. + * + * @api public + */ + +Socket.protocol = parser.protocol; // this is an int + +/** + * Expose deps for legacy compatibility + * and standalone browser access. + */ + +Socket.Socket = Socket; +Socket.Transport = _dereq_('./transport'); +Socket.transports = _dereq_('./transports'); +Socket.parser = _dereq_('engine.io-parser'); + +/** + * Creates transport of the given type. + * + * @param {String} transport name + * @return {Transport} + * @api private + */ + +Socket.prototype.createTransport = function (name) { + debug('creating transport "%s"', name); + var query = clone(this.query); + + // append engine.io protocol identifier + query.EIO = parser.protocol; + + // transport name + query.transport = name; + + // session id if we already have one + if (this.id) query.sid = this.id; + + var transport = new transports[name]({ + agent: this.agent, + hostname: this.hostname, + port: this.port, + secure: this.secure, + path: this.path, + query: query, + forceJSONP: this.forceJSONP, + forceBase64: this.forceBase64, + timestampRequests: this.timestampRequests, + timestampParam: this.timestampParam, + policyPort: this.policyPort, + socket: this + }); + + return transport; +}; + +function clone (obj) { + var o = {}; + for (var i in obj) { + if (obj.hasOwnProperty(i)) { + o[i] = obj[i]; + } + } + return o; +} + +/** + * Initializes transport to use and starts probe. + * + * @api private + */ +Socket.prototype.open = function () { + var transport; + if (this.rememberUpgrade && Socket.priorWebsocketSuccess && this.transports.indexOf('websocket') != -1) { + transport = 'websocket'; + } else { + transport = this.transports[0]; + } + this.readyState = 'opening'; + var transport = this.createTransport(transport); + transport.open(); + this.setTransport(transport); +}; + +/** + * Sets the current transport. Disables the existing one (if any). + * + * @api private + */ + +Socket.prototype.setTransport = function(transport){ + debug('setting transport %s', transport.name); + var self = this; + + if (this.transport) { + debug('clearing existing transport %s', this.transport.name); + this.transport.removeAllListeners(); + } + + // set up transport + this.transport = transport; + + // set up transport listeners + transport + .on('drain', function(){ + self.onDrain(); + }) + .on('packet', function(packet){ + self.onPacket(packet); + }) + .on('error', function(e){ + self.onError(e); + }) + .on('close', function(){ + self.onClose('transport close'); + }); +}; + +/** + * Probes a transport. + * + * @param {String} transport name + * @api private + */ + +Socket.prototype.probe = function (name) { + debug('probing transport "%s"', name); + var transport = this.createTransport(name, { probe: 1 }) + , failed = false + , self = this; + + Socket.priorWebsocketSuccess = false; + + function onTransportOpen(){ + if (self.onlyBinaryUpgrades) { + var upgradeLosesBinary = !this.supportsBinary && self.transport.supportsBinary; + failed = failed || upgradeLosesBinary; + } + if (failed) return; + + debug('probe transport "%s" opened', name); + transport.send([{ type: 'ping', data: 'probe' }]); + transport.once('packet', function (msg) { + if (failed) return; + if ('pong' == msg.type && 'probe' == msg.data) { + debug('probe transport "%s" pong', name); + self.upgrading = true; + self.emit('upgrading', transport); + Socket.priorWebsocketSuccess = 'websocket' == transport.name; + + debug('pausing current transport "%s"', self.transport.name); + self.transport.pause(function () { + if (failed) return; + if ('closed' == self.readyState || 'closing' == self.readyState) { + return; + } + debug('changing transport and sending upgrade packet'); + + cleanup(); + + self.setTransport(transport); + transport.send([{ type: 'upgrade' }]); + self.emit('upgrade', transport); + self.upgrading = false; + self.flush(); + }); + } else { + debug('probe transport "%s" failed', name); + var err = new Error('probe error'); + err.transport = transport.name; + self.emit('upgradeError', err); + } + }); + } + + function freezeTransport() { + if (failed) return; + + // Any callback called by transport should be ignored since now + failed = true; + cleanup(); + transport.close(); + } + + //Handle any error that happens while probing + function onerror(err) { + var error = new Error('probe error: ' + err); + error.transport = name; + + freezeTransport(); + + debug('probe transport "%s" failed because of error: %s', name, err); + + self.emit('upgradeError', error); + } + + function onTransportClose(){ + onerror("transport closed"); + } + + //When the socket is closed while we're probing + function onclose(){ + onerror("socket closed"); + } + + //When the socket is upgraded while we're probing + function onupgrade(to){ + if (transport && to.name != transport.name) { + debug('"%s" works - aborting "%s"', to.name, transport.name); + freezeTransport(); + } + } + + //Remove all listeners on the transport and on self + function cleanup(){ + transport.removeListener('open', onTransportOpen); + transport.removeListener('error', onerror); + transport.removeListener('close', onTransportClose); + self.removeListener('close', onclose); + self.removeListener('upgrading', onupgrade); + } + + transport.once('open', onTransportOpen); + transport.once('error', onerror); + transport.once('close', onTransportClose); + + this.once('close', onclose); + this.once('upgrading', onupgrade); + + transport.open(); + +}; + +/** + * Called when connection is deemed open. + * + * @api public + */ + +Socket.prototype.onOpen = function () { + debug('socket open'); + this.readyState = 'open'; + Socket.priorWebsocketSuccess = 'websocket' == this.transport.name; + this.emit('open'); + this.flush(); + + // we check for `readyState` in case an `open` + // listener already closed the socket + if ('open' == this.readyState && this.upgrade && this.transport.pause) { + debug('starting upgrade probes'); + for (var i = 0, l = this.upgrades.length; i < l; i++) { + this.probe(this.upgrades[i]); + } + } +}; + +/** + * Handles a packet. + * + * @api private + */ + +Socket.prototype.onPacket = function (packet) { + if ('opening' == this.readyState || 'open' == this.readyState) { + debug('socket receive: type "%s", data "%s"', packet.type, packet.data); + + this.emit('packet', packet); + + // Socket is live - any packet counts + this.emit('heartbeat'); + + switch (packet.type) { + case 'open': + this.onHandshake(parsejson(packet.data)); + break; + + case 'pong': + this.setPing(); + break; + + case 'error': + var err = new Error('server error'); + err.code = packet.data; + this.emit('error', err); + break; + + case 'message': + this.emit('data', packet.data); + this.emit('message', packet.data); + break; + } + } else { + debug('packet received with socket readyState "%s"', this.readyState); + } +}; + +/** + * Called upon handshake completion. + * + * @param {Object} handshake obj + * @api private + */ + +Socket.prototype.onHandshake = function (data) { + this.emit('handshake', data); + this.id = data.sid; + this.transport.query.sid = data.sid; + this.upgrades = this.filterUpgrades(data.upgrades); + this.pingInterval = data.pingInterval; + this.pingTimeout = data.pingTimeout; + this.onOpen(); + // In case open handler closes socket + if ('closed' == this.readyState) return; + this.setPing(); + + // Prolong liveness of socket on heartbeat + this.removeListener('heartbeat', this.onHeartbeat); + this.on('heartbeat', this.onHeartbeat); +}; + +/** + * Resets ping timeout. + * + * @api private + */ + +Socket.prototype.onHeartbeat = function (timeout) { + clearTimeout(this.pingTimeoutTimer); + var self = this; + self.pingTimeoutTimer = setTimeout(function () { + if ('closed' == self.readyState) return; + self.onClose('ping timeout'); + }, timeout || (self.pingInterval + self.pingTimeout)); +}; + +/** + * Pings server every `this.pingInterval` and expects response + * within `this.pingTimeout` or closes connection. + * + * @api private + */ + +Socket.prototype.setPing = function () { + var self = this; + clearTimeout(self.pingIntervalTimer); + self.pingIntervalTimer = setTimeout(function () { + debug('writing ping packet - expecting pong within %sms', self.pingTimeout); + self.ping(); + self.onHeartbeat(self.pingTimeout); + }, self.pingInterval); +}; + +/** +* Sends a ping packet. +* +* @api public +*/ + +Socket.prototype.ping = function () { + this.sendPacket('ping'); +}; + +/** + * Called on `drain` event + * + * @api private + */ + +Socket.prototype.onDrain = function() { + for (var i = 0; i < this.prevBufferLen; i++) { + if (this.callbackBuffer[i]) { + this.callbackBuffer[i](); + } + } + + this.writeBuffer.splice(0, this.prevBufferLen); + this.callbackBuffer.splice(0, this.prevBufferLen); + + // setting prevBufferLen = 0 is very important + // for example, when upgrading, upgrade packet is sent over, + // and a nonzero prevBufferLen could cause problems on `drain` + this.prevBufferLen = 0; + + if (this.writeBuffer.length == 0) { + this.emit('drain'); + } else { + this.flush(); + } +}; + +/** + * Flush write buffers. + * + * @api private + */ + +Socket.prototype.flush = function () { + if ('closed' != this.readyState && this.transport.writable && + !this.upgrading && this.writeBuffer.length) { + debug('flushing %d packets in socket', this.writeBuffer.length); + this.transport.send(this.writeBuffer); + // keep track of current length of writeBuffer + // splice writeBuffer and callbackBuffer on `drain` + this.prevBufferLen = this.writeBuffer.length; + this.emit('flush'); + } +}; + +/** + * Sends a message. + * + * @param {String} message. + * @param {Function} callback function. + * @return {Socket} for chaining. + * @api public + */ + +Socket.prototype.write = +Socket.prototype.send = function (msg, fn) { + this.sendPacket('message', msg, fn); + return this; +}; + +/** + * Sends a packet. + * + * @param {String} packet type. + * @param {String} data. + * @param {Function} callback function. + * @api private + */ + +Socket.prototype.sendPacket = function (type, data, fn) { + var packet = { type: type, data: data }; + this.emit('packetCreate', packet); + this.writeBuffer.push(packet); + this.callbackBuffer.push(fn); + this.flush(); +}; + +/** + * Closes the connection. + * + * @api private + */ + +Socket.prototype.close = function () { + if ('opening' == this.readyState || 'open' == this.readyState) { + this.onClose('forced close'); + debug('socket closing - telling transport to close'); + this.transport.close(); + } + + return this; +}; + +/** + * Called upon transport error + * + * @api private + */ + +Socket.prototype.onError = function (err) { + debug('socket error %j', err); + Socket.priorWebsocketSuccess = false; + this.emit('error', err); + this.onClose('transport error', err); +}; + +/** + * Called upon transport close. + * + * @api private + */ + +Socket.prototype.onClose = function (reason, desc) { + if ('opening' == this.readyState || 'open' == this.readyState) { + debug('socket close with reason: "%s"', reason); + var self = this; + + // clear timers + clearTimeout(this.pingIntervalTimer); + clearTimeout(this.pingTimeoutTimer); + + // clean buffers in next tick, so developers can still + // grab the buffers on `close` event + setTimeout(function() { + self.writeBuffer = []; + self.callbackBuffer = []; + self.prevBufferLen = 0; + }, 0); + + // stop event from firing again for transport + this.transport.removeAllListeners('close'); + + // ensure transport won't stay open + this.transport.close(); + + // ignore further transport communication + this.transport.removeAllListeners(); + + // set ready state + this.readyState = 'closed'; + + // clear session id + this.id = null; + + // emit close event + this.emit('close', reason, desc); + } +}; + +/** + * Filters upgrades, returning only those matching client transports. + * + * @param {Array} server upgrades + * @api private + * + */ + +Socket.prototype.filterUpgrades = function (upgrades) { + var filteredUpgrades = []; + for (var i = 0, j = upgrades.length; i