Compare commits

...

24 Commits

Author SHA1 Message Date
James Coglan 57e74231cf Bump version to 0.5.1. 2014-12-18 02:19:26 +00:00
James Coglan 5cba268409 Throw an error if drivers are created with unrecognised options. 2014-12-17 22:18:06 +00:00
James Coglan e8992add23 Bump version to 0.5.0 and document how to add extensions. 2014-12-13 13:00:47 +00:00
James Coglan 23675e08ee Don't expose the internal frame structure of messages to extensions; they should not be able to depend on seeing the actual frame boundaries off the wire because that requires all extensions to leave them intact. 2014-12-06 22:08:37 +00:00
James Coglan 66f46330e7 Remove some blank lines. 2014-12-05 19:36:06 +00:00
James Coglan 7376f48d6a Catch errors raised by Extensions.activate(). 2014-12-02 22:53:16 +00:00
James Coglan 330fa2073d Don't assume that messages coming out of extensions have anything but the standard websocket-extensions API. 2014-11-29 01:54:20 +00:00
James Coglan 58474a837c Populate the Message.data field before handing off to extensions. 2014-11-29 01:13:18 +00:00
James Coglan 82182eb348 Close the Extensions object when ending the session. 2014-11-29 00:55:02 +00:00
James Coglan 94041e7a86 websocket-extensions uses Node-style callbacks rather than emitting error events. 2014-11-27 09:32:49 +00:00
James Coglan e73f7db754 Rename validFrame() to validFrameRsv(). 2014-11-26 09:59:08 +00:00
James Coglan 0efc0ebb8d Don't pass the driver into the Extensions constructor. 2014-11-25 23:17:39 +00:00
James Coglan f904131d21 Extract the Extensions class into another module. 2014-11-25 21:14:17 +00:00
James Coglan ee664ba6c9 Decouple Extensions and individual extension libs from the Driver.fail() method, emitting errors instead. 2014-11-25 19:42:54 +00:00
James Coglan 538239b761 Move the header-parsing functions into class methods on Extensions. 2014-11-25 09:11:58 +00:00
James Coglan 2bc103c781 Adjust the data structures in the Extensions parser so that the server sets the order of extensions, per the spec. 2014-11-24 23:03:50 +00:00
James Coglan eceb59e2fb Implement the client-side framework for extensions. 2014-11-24 22:05:31 +00:00
James Coglan 40341c5439 The Client driver requires the crypto module for generating keys. 2014-11-24 00:41:34 +00:00
James Coglan 396637d463 Extract permessage-deflate into a separate module. 2014-11-23 21:46:04 +00:00
James Coglan c986ad231c Enable strict mode on all source files. 2014-11-23 20:09:26 +00:00
James Coglan ef897eed1c Implement deflate compression on server-to-client messages. 2014-11-23 19:45:02 +00:00
James Coglan 724a5f19a4 Flesh out the parameter validation for extensions, and provide better support for multiple values of extension params and HTTP headers. 2014-11-23 15:21:11 +00:00
James Coglan 0332ce2625 Initial support for permessage-deflate, only supports client-to-server compression but enough to pass the Autobahn suite. 2014-11-23 13:11:24 +00:00
James Coglan 6146615801 Extract Hybi parser state into Frame and Message structs.
This makes the structure of all the state in the parser easier to see,
it makes it easier to free memory by deleting a single struct rather
than many related fields, and it produces some data structures that will
be useful for adding extensions.
2014-11-22 13:34:34 +00:00
18 changed files with 358 additions and 171 deletions
+8
View File
@@ -1,3 +1,11 @@
### 0.5.1 / 2014-12-18
* Don't allow drivers to be created with unrecognized options
### 0.5.0 / 2014-12-13
* Support protocol extensions via the websocket-extensions module
### 0.4.0 / 2014-11-08
* Support connection via HTTP proxies using `CONNECT`
+10
View File
@@ -14,6 +14,9 @@ hook this module up to some I/O object, it will do all of this for you:
* Generate and send both server- and client-side handshakes
* Recognize when the handshake phase completes and the WS protocol begins
* Negotiate subprotocol selection based on `Sec-WebSocket-Protocol`
* Negotiate and use extensions via the
[websocket-extensions](https://github.com/faye/websocket-extensions-node)
module
* Buffer sent messages until the handshake process is finished
* Deal with proxies that defer delivery of the draft-76 handshake body
* Notify you when the socket is open and closed and when messages arrive
@@ -278,6 +281,13 @@ describing the error.
Sets the callback to execute when the socket becomes closed. The `event` object
has `code` and `reason` attributes.
#### `driver.addExtension(extension)`
Registers a protocol extension whose operation will be negotiated via the
`Sec-WebSocket-Extensions` header. `extension` is any extension compatible with
the [websocket-extensions](https://github.com/faye/websocket-extensions-node)
framework.
#### `driver.setHeader(name, value)`
Sets a custom header to be sent as part of the handshake response, either from
+3 -1
View File
@@ -1,8 +1,10 @@
var net = require('net'),
websocket = require('../lib/websocket/driver');
websocket = require('../lib/websocket/driver'),
deflate = require('permessage-deflate');
var server = net.createServer(function(connection) {
var driver = websocket.server();
driver.addExtension(deflate);
driver.on('connect', function() {
if (websocket.isWebSocket(driver)) driver.start();
+9 -2
View File
@@ -1,10 +1,13 @@
'use strict';
// Protocol references:
//
// * http://tools.ietf.org/html/draft-hixie-thewebsocketprotocol-75
// * http://tools.ietf.org/html/draft-hixie-thewebsocketprotocol-76
// * http://tools.ietf.org/html/draft-ietf-hybi-thewebsocketprotocol-17
var Client = require('./driver/client'),
var Base = require('./driver/base'),
Client = require('./driver/client'),
Server = require('./driver/server');
var Driver = {
@@ -35,8 +38,12 @@ var Driver = {
upgrade = request.headers.upgrade || '';
return request.method === 'GET' &&
connection.toLowerCase().split(/\s*,\s*/).indexOf('upgrade') >= 0 &&
connection.toLowerCase().split(/ *, */).indexOf('upgrade') >= 0 &&
upgrade.toLowerCase() === 'websocket';
},
validateOptions: function(options, validKeys) {
Base.validateOptions(options, validKeys);
}
};
+14
View File
@@ -1,3 +1,5 @@
'use strict';
var Emitter = require('events').EventEmitter,
util = require('util'),
streams = require('../streams'),
@@ -5,6 +7,7 @@ var Emitter = require('events').EventEmitter,
var Base = function(request, url, options) {
Emitter.call(this);
Base.validateOptions(options || {}, ['maxLength', 'masking', 'requireMasking', 'protocols']);
this._request = request;
this._options = options || {};
@@ -20,6 +23,13 @@ var Base = function(request, url, options) {
};
util.inherits(Base, Emitter);
Base.validateOptions = function(options, validKeys) {
for (var key in options) {
if (validKeys.indexOf(key) < 0)
throw new Error('Unrecognized option: ' + key);
}
};
var instance = {
// This is 64MB, small enough for an average VPS to handle without
// crashing from process out of memory
@@ -55,6 +65,10 @@ var instance = {
return this.STATES[this.readyState] || null;
},
addExtension: function(extension) {
return false;
},
setHeader: function(name, value) {
if (this.readyState > 0) return false;
this._headers.set(name, value);
+17 -5
View File
@@ -1,4 +1,7 @@
var url = require('url'),
'use strict';
var crypto = require('crypto'),
url = require('url'),
util = require('util'),
HttpParser = require('../http_parser'),
Base = require('./base'),
@@ -34,9 +37,7 @@ var Client = function(_url, options) {
util.inherits(Client, Hybi);
Client.generateKey = function() {
var buffer = new Buffer(16), i = buffer.length;
while (i--) buffer[i] = Math.floor(Math.random() * 256);
return buffer.toString('base64');
return crypto.randomBytes(16).toString('base64');
};
var instance = {
@@ -58,10 +59,17 @@ var instance = {
if (!this._http.isComplete()) return;
this._validateHandshake();
if (this.readyState === 3) return;
this._open();
this.parse(this._http.body);
},
_handshakeRequest: function() {
var extensions = this._extensions.generateOffer();
if (extensions)
this._headers.set('Sec-WebSocket-Extensions', extensions);
var start = 'GET ' + this._pathname + ' HTTP/1.1',
headers = [start, this._headers.toString(), ''];
@@ -110,7 +118,11 @@ var instance = {
this.protocol = protocol;
}
this._open();
try {
this._extensions.activate(this.headers['sec-websocket-extensions']);
} catch (e) {
return this._failHandshake(e.message);
}
}
};
+2
View File
@@ -1,3 +1,5 @@
'use strict';
var Base = require('./base'),
util = require('util');
+2 -3
View File
@@ -1,3 +1,5 @@
'use strict';
var Base = require('./base'),
Draft75 = require('./draft75'),
crypto = require('crypto'),
@@ -66,13 +68,10 @@ var instance = {
var body = new Buffer(this._body.slice(0, this.BODY_SIZE));
var headers = this._request.headers,
key1 = headers['sec-websocket-key1'],
value1 = numberFromKey(key1) / spacesInKey(key1),
key2 = headers['sec-websocket-key2'],
value2 = numberFromKey(key2) / spacesInKey(key2),
md5 = crypto.createHash('md5');
md5.update(bigEndian(value1));
+2
View File
@@ -1,3 +1,5 @@
'use strict';
var Headers = function() {
this.clear();
};
+198 -151
View File
@@ -1,12 +1,17 @@
var crypto = require('crypto'),
util = require('util'),
Base = require('./base'),
Reader = require('./hybi/stream_reader');
'use strict';
var crypto = require('crypto'),
util = require('util'),
Extensions = require('websocket-extensions'),
Base = require('./base'),
Frame = require('./hybi/frame'),
Message = require('./hybi/message'),
Reader = require('./hybi/stream_reader');
var Hybi = function(request, url, options) {
Base.apply(this, arguments);
this._reset();
this._extensions = new Extensions();
this._reader = new Reader();
this._stage = 0;
this._masking = this._options.masking;
@@ -15,7 +20,7 @@ var Hybi = function(request, url, options) {
this._pingCallbacks = {};
if (typeof this._protocols === 'string')
this._protocols = this._protocols.split(/\s*,\s*/);
this._protocols = this._protocols.split(/ *, */);
if (!this._request) return;
@@ -29,7 +34,7 @@ var Hybi = function(request, url, options) {
this._headers.set('Sec-WebSocket-Accept', Hybi.generateAccept(secKey));
if (protos !== undefined) {
if (typeof protos === 'string') protos = protos.split(/\s*,\s*/);
if (typeof protos === 'string') protos = protos.split(/ *, */);
this.protocol = protos.filter(function(p) { return supported.indexOf(p) >= 0 })[0];
if (this.protocol) this._headers.set('Sec-WebSocket-Protocol', this.protocol);
}
@@ -75,9 +80,9 @@ var instance = {
pong: 10
},
OPCODE_CODES: [0, 1, 2, 8, 9, 10],
FRAGMENTED_OPCODES: [0, 1, 2],
OPENING_OPCODES: [1, 2],
OPCODE_CODES: [0, 1, 2, 8, 9, 10],
MESSAGE_OPCODES: [0, 1, 2],
OPENING_OPCODES: [1, 2],
TWO_POWERS: [0, 1, 2, 3, 4, 5, 6, 7].map(function(n) { return Math.pow(2, 8 * n) }),
@@ -100,6 +105,11 @@ var instance = {
// http://www.w3.org/International/questions/qa-forms-utf-8.en.php
UTF8_MATCH: /^([\x00-\x7F]|[\xC2-\xDF][\x80-\xBF]|\xE0[\xA0-\xBF][\x80-\xBF]|[\xE1-\xEC\xEE\xEF][\x80-\xBF]{2}|\xED[\x80-\x9F][\x80-\xBF]|\xF0[\x90-\xBF][\x80-\xBF]{2}|[\xF1-\xF3][\x80-\xBF]{3}|\xF4[\x80-\x8F][\x80-\xBF]{2})*$/,
addExtension: function(extension) {
this._extensions.add(extension);
return true;
},
parse: function(data) {
this._reader.put(data);
var buffer = true;
@@ -116,20 +126,20 @@ var instance = {
break;
case 2:
buffer = this._reader.read(this._lengthSize);
buffer = this._reader.read(this._frame.lengthBytes);
if (buffer) this._parseExtendedLength(buffer);
break;
case 3:
buffer = this._reader.read(4);
if (buffer) {
this._mask = buffer;
this._frame.maskingKey = buffer;
this._stage = 4;
}
break;
case 4:
buffer = this._reader.read(this._length);
buffer = this._reader.read(this._frame.length);
if (buffer) {
this._emitFrame(buffer);
this._stage = 0;
@@ -142,61 +152,6 @@ var instance = {
}
},
frame: function(data, type, code) {
if (this.readyState <= 0) return this._queue([data, type, code]);
if (this.readyState !== 1) return false;
if (data instanceof Array) data = new Buffer(data);
var isText = (typeof data === 'string'),
opcode = this.OPCODES[type || (isText ? 'text' : 'binary')],
buffer = isText ? new Buffer(data, 'utf8') : data,
insert = code ? 2 : 0,
length = buffer.length + insert,
header = (length <= 125) ? 2 : (length <= 65535 ? 4 : 10),
offset = header + (this._masking ? 4 : 0),
masked = this._masking ? this.MASK : 0,
frame = new Buffer(length + offset),
BYTE = this.BYTE,
mask, i;
frame[0] = this.FIN | opcode;
if (length <= 125) {
frame[1] = masked | length;
} else if (length <= 65535) {
frame[1] = masked | 126;
frame[2] = Math.floor(length / 256);
frame[3] = length & BYTE;
} else {
frame[1] = masked | 127;
frame[2] = Math.floor(length / Math.pow(2,56)) & BYTE;
frame[3] = Math.floor(length / Math.pow(2,48)) & BYTE;
frame[4] = Math.floor(length / Math.pow(2,40)) & BYTE;
frame[5] = Math.floor(length / Math.pow(2,32)) & BYTE;
frame[6] = Math.floor(length / Math.pow(2,24)) & BYTE;
frame[7] = Math.floor(length / Math.pow(2,16)) & BYTE;
frame[8] = Math.floor(length / Math.pow(2,8)) & BYTE;
frame[9] = length & BYTE;
}
if (code) {
frame[offset] = Math.floor(code / 256) & BYTE;
frame[offset+1] = code & BYTE;
}
buffer.copy(frame, offset + insert);
if (this._masking) {
mask = [Math.floor(Math.random() * 256), Math.floor(Math.random() * 256),
Math.floor(Math.random() * 256), Math.floor(Math.random() * 256)];
new Buffer(mask).copy(frame, header);
Hybi.mask(frame, mask, offset);
}
this._write(frame);
return true;
},
text: function(message) {
return this.frame(message, 'text');
},
@@ -228,7 +183,104 @@ var instance = {
}
},
frame: function(data, type, code) {
if (this.readyState <= 0) return this._queue([data, type, code]);
if (this.readyState !== 1) return false;
if (data instanceof Array) data = new Buffer(data);
var message = new Message(),
isText = (typeof data === 'string'),
payload, buffer;
message.rsv1 = message.rsv2 = message.rsv3 = false;
message.opcode = this.OPCODES[type || (isText ? 'text' : 'binary')];
payload = isText ? new Buffer(data, 'utf8') : data;
if (code) {
buffer = payload;
payload = new Buffer(2 + buffer.length);
payload[0] = ~~(code / 256) & this.BYTE;
payload[1] = code & this.BYTE;
buffer.copy(payload, 2);
}
message.data = payload;
var onMessageReady = function(message) {
var frame = new Frame();
frame.final = true;
frame.rsv1 = message.rsv1;
frame.rsv2 = message.rsv2;
frame.rsv3 = message.rsv3;
frame.opcode = message.opcode;
frame.masked = !!this._masking;
frame.length = message.data.length;
frame.payload = message.data;
if (frame.masked) frame.maskingKey = crypto.randomBytes(4);
this._sendFrame(frame);
};
if (this.MESSAGE_OPCODES.indexOf(message.opcode) >= 0)
this._extensions.processOutgoingMessage(message, function(error, message) {
if (error) return this._fail('extension_error', error.message);
onMessageReady.call(this, message);
}, this);
else
onMessageReady.call(this, message);
return true;
},
_sendFrame: function(frame) {
var length = frame.length,
header = (length <= 125) ? 2 : (length <= 65535 ? 4 : 10),
offset = header + (frame.masked ? 4 : 0),
buffer = new Buffer(offset + length),
BYTE = this.BYTE,
masked = frame.masked ? this.MASK : 0;
buffer[0] = (frame.final ? this.FIN : 0) |
(frame.rsv1 ? this.RSV1 : 0) |
(frame.rsv2 ? this.RSV2 : 0) |
(frame.rsv3 ? this.RSV3 : 0) |
frame.opcode;
if (length <= 125) {
buffer[1] = masked | length;
} else if (length <= 65535) {
buffer[1] = masked | 126;
buffer[2] = ~~(length / 256);
buffer[3] = length & BYTE;
} else {
buffer[1] = masked | 127;
buffer[2] = ~~(length / Math.pow(2, 56)) & BYTE;
buffer[3] = ~~(length / Math.pow(2, 48)) & BYTE;
buffer[4] = ~~(length / Math.pow(2, 40)) & BYTE;
buffer[5] = ~~(length / Math.pow(2, 32)) & BYTE;
buffer[6] = ~~(length / Math.pow(2, 24)) & BYTE;
buffer[7] = ~~(length / Math.pow(2, 16)) & BYTE;
buffer[8] = ~~(length / Math.pow(2, 8)) & BYTE;
buffer[9] = length & BYTE;
}
if (frame.masked) {
frame.maskingKey.copy(buffer, header);
Hybi.mask(frame.payload, frame.maskingKey).copy(buffer, offset);
} else {
frame.payload.copy(buffer, offset);
}
this._write(buffer);
},
_handshakeResponse: function() {
var extensions = this._extensions.generateResponse(this._request.headers['sec-websocket-extensions']);
if (extensions) this._headers.set('Sec-WebSocket-Extensions', extensions);
var start = 'HTTP/1.1 101 Switching Protocols',
headers = [start, this._headers.toString(), ''];
@@ -237,9 +289,12 @@ var instance = {
_shutdown: function(code, reason) {
this.frame(reason, 'close', code);
delete this._frame;
delete this._message;
this.readyState = 3;
this._stage = 5;
this.emit('close', new Base.CloseEvent(code, reason));
this._extensions.close();
},
_fail: function(type, message) {
@@ -252,56 +307,66 @@ var instance = {
return (data & rsv) === rsv;
});
if (rsvs.filter(function(rsv) { return rsv }).length > 0)
var frame = this._frame = new Frame();
frame.final = (data & this.FIN) === this.FIN;
frame.rsv1 = rsvs[0];
frame.rsv2 = rsvs[1];
frame.rsv3 = rsvs[2];
frame.opcode = (data & this.OPCODE);
if (!this._extensions.validFrameRsv(frame))
return this._fail('protocol_error',
'One or more reserved bits are on: reserved1 = ' + (rsvs[0] ? 1 : 0) +
', reserved2 = ' + (rsvs[1] ? 1 : 0) +
', reserved3 = ' + (rsvs[2] ? 1 : 0));
'One or more reserved bits are on: reserved1 = ' + (frame.rsv1 ? 1 : 0) +
', reserved2 = ' + (frame.rsv2 ? 1 : 0) +
', reserved3 = ' + (frame.rsv3 ? 1 : 0));
this._final = (data & this.FIN) === this.FIN;
this._opcode = (data & this.OPCODE);
if (this.OPCODE_CODES.indexOf(frame.opcode) < 0)
return this._fail('protocol_error', 'Unrecognized frame opcode: ' + frame.opcode);
if (this.OPCODE_CODES.indexOf(this._opcode) < 0)
return this._fail('protocol_error', 'Unrecognized frame opcode: ' + this._opcode);
if (this.MESSAGE_OPCODES.indexOf(frame.opcode) < 0 && !frame.final)
return this._fail('protocol_error', 'Received fragmented control frame: opcode = ' + frame.opcode);
if (this.FRAGMENTED_OPCODES.indexOf(this._opcode) < 0 && !this._final)
return this._fail('protocol_error', 'Received fragmented control frame: opcode = ' + this._opcode);
if (this._mode && this.OPENING_OPCODES.indexOf(this._opcode) >= 0)
if (this._message && this.OPENING_OPCODES.indexOf(frame.opcode) >= 0)
return this._fail('protocol_error', 'Received new data frame but previous continuous frame is unfinished');
this._stage = 1;
},
_parseLength: function(data) {
this._masked = (data & this.MASK) === this.MASK;
if (this._requireMasking && !this._masked)
var frame = this._frame;
frame.masked = (data & this.MASK) === this.MASK;
if (this._requireMasking && !frame.masked)
return this._fail('unacceptable', 'Received unmasked frame but masking is required');
this._length = (data & this.LENGTH);
frame.length = (data & this.LENGTH);
if (this._length >= 0 && this._length <= 125) {
if (frame.length >= 0 && frame.length <= 125) {
if (!this._checkFrameLength()) return;
this._stage = this._masked ? 3 : 4;
this._stage = frame.masked ? 3 : 4;
} else {
this._lengthSize = (this._length === 126 ? 2 : 8);
this._stage = 2;
frame.lengthBytes = (frame.length === 126 ? 2 : 8);
this._stage = 2;
}
},
_parseExtendedLength: function(buffer) {
this._length = this._getInteger(buffer);
var frame = this._frame;
frame.length = this._getInteger(buffer);
if (this.FRAGMENTED_OPCODES.indexOf(this._opcode) < 0 && this._length > 125)
return this._fail('protocol_error', 'Received control frame having too long payload: ' + this._length);
if (this.MESSAGE_OPCODES.indexOf(frame.opcode) < 0 && frame.length > 125)
return this._fail('protocol_error', 'Received control frame having too long payload: ' + frame.length);
if (!this._checkFrameLength()) return;
this._stage = this._masked ? 3 : 4;
this._stage = frame.masked ? 3 : 4;
},
_checkFrameLength: function() {
if (this.__blength + this._length > this._maxLength) {
var length = this._message ? this._message.length : 0;
if (length + this._frame.length > this._maxLength) {
this._fail('too_large', 'WebSocket frame length too large');
return false;
} else {
@@ -310,48 +375,31 @@ var instance = {
},
_emitFrame: function(buffer) {
var payload = Hybi.mask(buffer, this._mask),
isFinal = this._final,
opcode = this._opcode;
var frame = this._frame,
payload = frame.payload = Hybi.mask(buffer, frame.maskingKey),
opcode = frame.opcode,
message,
code, reason,
callbacks, callback;
this._final = this._opcode = this._length = this._lengthSize = this._masked = this._mask = null;
delete this._frame;
if (opcode === this.OPCODES.continuation) {
if (!this._mode) return this._fail('protocol_error', 'Received unexpected continuation frame');
this._buffer(payload);
if (isFinal) {
var message = this._concatBuffer();
if (this._mode === 'text') message = this._encode(message);
this._reset();
if (message === null)
this._fail('encoding_error', 'Could not decode a text frame as UTF-8');
else
this.emit('message', new Base.MessageEvent(message));
}
if (!this._message) return this._fail('protocol_error', 'Received unexpected continuation frame');
this._message.pushFrame(frame);
}
else if (opcode === this.OPCODES.text) {
if (isFinal) {
var message = this._encode(payload);
if (message === null)
this._fail('encoding_error', 'Could not decode a text frame as UTF-8');
else
this.emit('message', new Base.MessageEvent(message));
} else {
this._mode = 'text';
this._buffer(payload);
}
if (opcode === this.OPCODES.text || opcode === this.OPCODES.binary) {
this._message = new Message();
this._message.pushFrame(frame);
}
else if (opcode === this.OPCODES.binary) {
if (isFinal) {
this.emit('message', new Base.MessageEvent(payload));
} else {
this._mode = 'binary';
this._buffer(payload);
}
}
else if (opcode === this.OPCODES.close) {
var code = (payload.length >= 2) ? 256 * payload[0] + payload[1] : null,
reason = (payload.length > 2) ? this._encode(payload.slice(2)) : null;
if (frame.final && this.MESSAGE_OPCODES.indexOf(opcode) >= 0)
return this._emitMessage(this._message);
if (opcode === this.OPCODES.close) {
code = (payload.length >= 2) ? 256 * payload[0] + payload[1] : null;
reason = (payload.length > 2) ? this._encode(payload.slice(2)) : null;
if (!(payload.length === 0) &&
!(code !== null && code >= this.MIN_RESERVED_ERROR && code <= this.MAX_RESERVED_ERROR) &&
@@ -363,39 +411,38 @@ var instance = {
this._shutdown(code, reason || '');
}
else if (opcode === this.OPCODES.ping) {
if (opcode === this.OPCODES.ping) {
this.frame(payload, 'pong');
}
else if (opcode === this.OPCODES.pong) {
var callbacks = this._pingCallbacks,
message = this._encode(payload),
callback = callbacks[message];
if (opcode === this.OPCODES.pong) {
callbacks = this._pingCallbacks;
message = this._encode(payload);
callback = callbacks[message];
delete callbacks[message];
if (callback) callback()
}
},
_buffer: function(fragment) {
this.__buffer.push(fragment);
this.__blength += fragment.length;
},
_emitMessage: function(message) {
var message = this._message;
message.read();
_concatBuffer: function() {
var buffer = new Buffer(this.__blength),
offset = 0;
delete this._message;
for (var i = 0, n = this.__buffer.length; i < n; i++) {
this.__buffer[i].copy(buffer, offset);
offset += this.__buffer[i].length;
}
return buffer;
},
this._extensions.processIncomingMessage(message, function(error, message) {
if (error) return this._fail('extension_error', error.message);
_reset: function() {
this._mode = null;
this.__buffer = [];
this.__blength = 0;
var payload = message.data;
if (message.opcode === this.OPCODES.text) payload = this._encode(payload);
if (payload === null)
return this._fail('encoding_error', 'Could not decode a text frame as UTF-8');
else
this.emit('message', new Base.MessageEvent(payload));
}, this);
},
_encode: function(buffer) {
+21
View File
@@ -0,0 +1,21 @@
'use strict';
var Frame = function() {};
var instance = {
final: false,
rsv1: false,
rsv2: false,
rsv3: false,
opcode: null,
masked: false,
maskingKey: null,
lengthBytes: 1,
length: 0,
payload: null
};
for (var key in instance)
Frame.prototype[key] = instance[key];
module.exports = Frame;
+41
View File
@@ -0,0 +1,41 @@
'use strict';
var Message = function() {
this.rsv1 = false;
this.rsv2 = false;
this.rsv3 = false;
this.opcode = null
this.length = 0;
this._chunks = [];
};
var instance = {
read: function() {
if (this.data) return this.data;
this.data = new Buffer(this.length);
var offset = 0;
for (var i = 0, n = this._chunks.length; i < n; i++) {
this._chunks[i].copy(this.data, offset);
offset += this._chunks[i].length;
}
return this.data;
},
pushFrame: function(frame) {
this.rsv1 = this.rsv1 || frame.rsv1;
this.rsv2 = this.rsv2 || frame.rsv2;
this.rsv3 = this.rsv3 || frame.rsv3;
if (this.opcode === null) this.opcode = frame.opcode;
this._chunks.push(frame.payload);
this.length += frame.length;
}
};
for (var key in instance)
Message.prototype[key] = instance[key];
module.exports = Message;
@@ -1,3 +1,5 @@
'use strict';
var StreamReader = function() {
this._queue = [];
this._queueSize = 0;
+2
View File
@@ -1,3 +1,5 @@
'use strict';
var Stream = require('stream').Stream,
url = require('url'),
util = require('util'),
+4 -2
View File
@@ -1,3 +1,5 @@
'use strict';
var util = require('util'),
HttpParser = require('../http_parser'),
Base = require('./base'),
@@ -34,8 +36,8 @@ var instance = {
this._delegate = Server.http(this, this._options);
this._delegate.messages = this.messages;
this._delegate.io = this.io;
this._open();
this._delegate.on('open', function() { self._open() });
this.EVENTS.forEach(function(event) {
this._delegate.on(event, function(e) { self.emit(event, e) });
}, this);
@@ -55,7 +57,7 @@ var instance = {
}
};
['setHeader', 'start', 'state', 'frame', 'text', 'binary', 'ping', 'close'].forEach(function(method) {
['addExtension', 'setHeader', 'start', 'frame', 'text', 'binary', 'ping', 'close'].forEach(function(method) {
instance[method] = function() {
if (this._delegate) {
return this._delegate[method].apply(this._delegate, arguments);
+17 -4
View File
@@ -1,3 +1,5 @@
'use strict';
var HTTPParser = process.binding('http_parser').HTTPParser,
version = HTTPParser.RESPONSE ? 6 : 4;
@@ -19,7 +21,12 @@ var HttpParser = function(type) {
};
this._parser.onHeaderValue = function(b, start, length) {
self.headers[current] = b.toString('utf8', start, start + length);
var value = b.toString('utf8', start, start + length);
if (self.headers.hasOwnProperty(current))
self.headers[current] += ', ' + value;
else
self.headers[current] = value;
};
this._parser.onHeadersComplete = this._parser[HTTPParser.kOnHeadersComplete] = function(info) {
@@ -27,11 +34,17 @@ var HttpParser = function(type) {
self.statusCode = info.statusCode;
self.url = info.url;
var headers = info.headers;
var headers = info.headers, key, value;
if (!headers) return;
for (var i = 0, n = headers.length; i < n; i += 2)
self.headers[headers[i].toLowerCase()] = headers[i+1];
for (var i = 0, n = headers.length; i < n; i += 2) {
key = headers[i].toLowerCase();
value = headers[i+1];
if (self.headers.hasOwnProperty(key))
self.headers[key] += ', ' + value;
else
self.headers[key] = value;
}
self._complete = true;
};
+2
View File
@@ -1,3 +1,5 @@
'use strict';
/**
Streams in a WebSocket connection
+4 -3
View File
@@ -5,10 +5,11 @@
, "keywords" : ["websocket"]
, "license" : "MIT"
, "version" : "0.4.0"
, "engines" : {"node": ">=0.4.0"}
, "version" : "0.5.1"
, "engines" : {"node": ">=0.6.0"}
, "main" : "./lib/websocket/driver"
, "devDependencies" : {"jstest": ""}
, "dependencies" : {"websocket-extensions": ">=0.1.0"}
, "devDependencies" : {"jstest": "", "permessage-deflate": ""}
, "scripts" : {"test": "jstest spec/runner.js"}