Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c84cce9a1f | |||
| c8e5f1eb78 | |||
| eb5e39e7e3 | |||
| a95f68080e | |||
| 7a1cb15b8c | |||
| 14937dd365 | |||
| 780dd331fe | |||
| 4a3d51cc5b | |||
| ed17186354 | |||
| 1df1293456 | |||
| d901d3e48d | |||
| a49cd60cb7 | |||
| 4dd872f0ef | |||
| 90b8c9d23a | |||
| c3ac50931e | |||
| 5f6873ebc0 | |||
| 309b5651a7 |
@@ -1,3 +1,13 @@
|
||||
### 0.6.1 / 2015-07-13
|
||||
|
||||
* Use the `buffer.{read,write}UInt{16,32}BE` methods for reading/writing numbers
|
||||
to buffers rather than including duplicate logic for this
|
||||
|
||||
### 0.6.0 / 2015-07-08
|
||||
|
||||
* Allow the parser to recover cleanly if event listeners raise an error
|
||||
* Add a `pong` method for sending unsolicited pong frames
|
||||
|
||||
### 0.5.4 / 2015-03-29
|
||||
|
||||
* Don't emit extra close frames if we receive a close frame after we already
|
||||
|
||||
@@ -21,7 +21,7 @@ this module up to some I/O object, it will do all of this for you:
|
||||
* 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
|
||||
* Recombine fragmented messages
|
||||
* Dispatch text, binary, ping and close frames
|
||||
* Dispatch text, binary, ping, pong and close frames
|
||||
* Manage the socket-closing handshake process
|
||||
* Automatically reply to ping frames with a matching pong
|
||||
* Apply masking to messages sent by the client
|
||||
@@ -259,11 +259,11 @@ the `driver.io` stream.
|
||||
|
||||
#### `driver.on('open', function(event) {})`
|
||||
|
||||
Sets the callback to execute when the socket becomes open.
|
||||
Adds a callback to execute when the socket becomes open.
|
||||
|
||||
#### `driver.on('message', function(event) {})`
|
||||
|
||||
Sets the callback to execute when a message is received. `event` will have a
|
||||
Adds a callback to execute when a message is received. `event` will have a
|
||||
`data` attribute containing either a string in the case of a text message or a
|
||||
`Buffer` in the case of a binary message.
|
||||
|
||||
@@ -272,13 +272,13 @@ which emits strings for text messages and buffers for binary messages.
|
||||
|
||||
#### `driver.on('error', function(event) {})`
|
||||
|
||||
Sets the callback to execute when a protocol error occurs due to the other peer
|
||||
Adds a callback to execute when a protocol error occurs due to the other peer
|
||||
sending an invalid byte sequence. `event` will have a `message` attribute
|
||||
describing the error.
|
||||
|
||||
#### `driver.on('close', function(event) {})`
|
||||
|
||||
Sets the callback to execute when the socket becomes closed. The `event` object
|
||||
Adds a callback to execute when the socket becomes closed. The `event` object
|
||||
has `code` and `reason` attributes.
|
||||
|
||||
#### `driver.addExtension(extension)`
|
||||
@@ -330,6 +330,15 @@ callback are both optional. If a callback is given, it will be invoked when the
|
||||
socket receives a pong frame whose content matches `string`. Returns `false` if
|
||||
frames can no longer be sent, or if the driver does not support ping/pong.
|
||||
|
||||
#### `driver.pong(string = '')`
|
||||
|
||||
Sends a pong frame over the socket, queueing it if necessary. `string` is
|
||||
optional. Returns `false` if frames can no longer be sent, or if the driver does
|
||||
not support ping/pong.
|
||||
|
||||
You don't need to call this when a ping frame is received; pings are replied to
|
||||
automatically by the driver. This method is for sending unsolicited pongs.
|
||||
|
||||
#### `driver.close()`
|
||||
|
||||
Initiates the closing handshake if the socket is still open. For drivers with no
|
||||
|
||||
@@ -3,13 +3,15 @@
|
||||
var Emitter = require('events').EventEmitter,
|
||||
util = require('util'),
|
||||
streams = require('../streams'),
|
||||
Headers = require('./headers');
|
||||
Headers = require('./headers'),
|
||||
Reader = require('./stream_reader');
|
||||
|
||||
var Base = function(request, url, options) {
|
||||
Emitter.call(this);
|
||||
Base.validateOptions(options || {}, ['maxLength', 'masking', 'requireMasking', 'protocols']);
|
||||
|
||||
this._request = request;
|
||||
this._reader = new Reader();
|
||||
this._options = options || {};
|
||||
this._maxLength = this._options.maxLength || this.MAX_LENGTH;
|
||||
this._headers = new Headers();
|
||||
@@ -96,6 +98,10 @@ var instance = {
|
||||
return false;
|
||||
},
|
||||
|
||||
pong: function() {
|
||||
return false;
|
||||
},
|
||||
|
||||
close: function(reason, code) {
|
||||
if (this.readyState !== 1) return false;
|
||||
this.readyState = 3;
|
||||
|
||||
@@ -52,11 +52,11 @@ var instance = {
|
||||
return true;
|
||||
},
|
||||
|
||||
parse: function(data) {
|
||||
parse: function(chunk) {
|
||||
if (this.readyState === 3) return;
|
||||
if (this.readyState > 0) return Hybi.prototype.parse.call(this, data);
|
||||
if (this.readyState > 0) return Hybi.prototype.parse.call(this, chunk);
|
||||
|
||||
this._http.parse(data);
|
||||
this._http.parse(chunk);
|
||||
if (!this._http.isComplete()) return;
|
||||
|
||||
this._validateHandshake();
|
||||
@@ -79,8 +79,8 @@ var instance = {
|
||||
|
||||
_failHandshake: function(message) {
|
||||
message = 'Error during WebSocket handshake: ' + message;
|
||||
this.emit('error', new Error(message));
|
||||
this.readyState = 3;
|
||||
this.emit('error', new Error(message));
|
||||
this.emit('close', new Base.CloseEvent(this.ERRORS.protocol_error, message));
|
||||
},
|
||||
|
||||
|
||||
@@ -5,7 +5,7 @@ var Base = require('./base'),
|
||||
|
||||
var Draft75 = function(request, url, options) {
|
||||
Base.apply(this, arguments);
|
||||
this._stage = 0;
|
||||
this._stage = 0;
|
||||
this.version = 'hixie-75';
|
||||
|
||||
this._headers.set('Upgrade', 'WebSocket');
|
||||
@@ -23,46 +23,46 @@ var instance = {
|
||||
return true;
|
||||
},
|
||||
|
||||
parse: function(buffer) {
|
||||
parse: function(chunk) {
|
||||
if (this.readyState > 1) return;
|
||||
|
||||
var data, message, value;
|
||||
for (var i = 0, n = buffer.length; i < n; i++) {
|
||||
data = buffer[i];
|
||||
this._reader.put(chunk);
|
||||
|
||||
this._reader.eachByte(function(octet) {
|
||||
var message;
|
||||
|
||||
switch (this._stage) {
|
||||
case -1:
|
||||
this._body.push(data);
|
||||
this._body.push(octet);
|
||||
this._sendHandshakeBody();
|
||||
break;
|
||||
|
||||
case 0:
|
||||
this._parseLeadingByte(data);
|
||||
this._parseLeadingByte(octet);
|
||||
break;
|
||||
|
||||
case 1:
|
||||
value = (data & 0x7F);
|
||||
this._length = value + 128 * this._length;
|
||||
this._length = (octet & 0x7F) + 128 * this._length;
|
||||
|
||||
if (this._closing && this._length === 0) {
|
||||
return this.close();
|
||||
}
|
||||
else if ((0x80 & data) !== 0x80) {
|
||||
else if ((octet & 0x80) !== 0x80) {
|
||||
if (this._length === 0) {
|
||||
this._stage = 0;
|
||||
}
|
||||
else {
|
||||
this._skipped = 0;
|
||||
this._stage = 2;
|
||||
this._stage = 2;
|
||||
}
|
||||
}
|
||||
break;
|
||||
|
||||
case 2:
|
||||
if (data === 0xFF) {
|
||||
if (octet === 0xFF) {
|
||||
this._stage = 0;
|
||||
message = new Buffer(this._buffer).toString('utf8', 0, this._buffer.length);
|
||||
this.emit('message', new Base.MessageEvent(message));
|
||||
this._stage = 0;
|
||||
}
|
||||
else {
|
||||
if (this._length) {
|
||||
@@ -70,25 +70,25 @@ var instance = {
|
||||
if (this._skipped === this._length)
|
||||
this._stage = 0;
|
||||
} else {
|
||||
this._buffer.push(data);
|
||||
this._buffer.push(octet);
|
||||
if (this._buffer.length > this._maxLength) return this.close();
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
}, this);
|
||||
},
|
||||
|
||||
frame: function(data) {
|
||||
if (this.readyState === 0) return this._queue([data]);
|
||||
frame: function(buffer) {
|
||||
if (this.readyState === 0) return this._queue([buffer]);
|
||||
if (this.readyState > 1) return false;
|
||||
|
||||
var buffer = new Buffer(data, 'utf8'),
|
||||
frame = new Buffer(buffer.length + 2);
|
||||
var payload = new Buffer(buffer, 'utf8'),
|
||||
frame = new Buffer(payload.length + 2);
|
||||
|
||||
frame[0] = 0x00;
|
||||
frame[buffer.length + 1] = 0xFF;
|
||||
buffer.copy(frame, 1);
|
||||
frame[payload.length + 1] = 0xFF;
|
||||
payload.copy(frame, 1);
|
||||
|
||||
this._write(frame);
|
||||
return true;
|
||||
@@ -101,15 +101,15 @@ var instance = {
|
||||
return new Buffer(headers.join('\r\n'), 'utf8');
|
||||
},
|
||||
|
||||
_parseLeadingByte: function(data) {
|
||||
if ((0x80 & data) === 0x80) {
|
||||
_parseLeadingByte: function(octet) {
|
||||
if ((octet & 0x80) === 0x80) {
|
||||
this._length = 0;
|
||||
this._stage = 1;
|
||||
this._stage = 1;
|
||||
} else {
|
||||
delete this._length;
|
||||
delete this._skipped;
|
||||
this._buffer = [];
|
||||
this._stage = 2;
|
||||
this._stage = 2;
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
@@ -14,14 +14,6 @@ var spacesInKey = function(key) {
|
||||
return key.match(/ /g).length;
|
||||
};
|
||||
|
||||
var bigEndian = function(number) {
|
||||
var string = '';
|
||||
[24, 16, 8, 0].forEach(function(offset) {
|
||||
string += String.fromCharCode(number >> offset & 0xFF);
|
||||
});
|
||||
return string;
|
||||
};
|
||||
|
||||
|
||||
var Draft76 = function(request, url, options) {
|
||||
Draft75.apply(this, arguments);
|
||||
@@ -65,19 +57,20 @@ var instance = {
|
||||
|
||||
_handshakeSignature: function() {
|
||||
if (this._body.length < this.BODY_SIZE) return null;
|
||||
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 = crypto.createHash('md5'),
|
||||
buffer = new Buffer(8 + this.BODY_SIZE);
|
||||
|
||||
md5.update(bigEndian(value1));
|
||||
md5.update(bigEndian(value2));
|
||||
md5.update(body.toString('binary'));
|
||||
buffer.writeUInt32BE(value1, 0);
|
||||
buffer.writeUInt32BE(value2, 4);
|
||||
new Buffer(this._body).copy(buffer, 8, 0, this.BODY_SIZE);
|
||||
|
||||
md5.update(buffer);
|
||||
return new Buffer(md5.digest('binary'), 'binary');
|
||||
},
|
||||
|
||||
@@ -94,9 +87,9 @@ var instance = {
|
||||
this.parse(this._body.slice(this.BODY_SIZE));
|
||||
},
|
||||
|
||||
_parseLeadingByte: function(data) {
|
||||
if (data !== 0xFF)
|
||||
return Draft75.prototype._parseLeadingByte.call(this, data);
|
||||
_parseLeadingByte: function(octet) {
|
||||
if (octet !== 0xFF)
|
||||
return Draft75.prototype._parseLeadingByte.call(this, octet);
|
||||
|
||||
this._closing = true;
|
||||
this._length = 0;
|
||||
|
||||
@@ -5,14 +5,12 @@ var crypto = require('crypto'),
|
||||
Extensions = require('websocket-extensions'),
|
||||
Base = require('./base'),
|
||||
Frame = require('./hybi/frame'),
|
||||
Message = require('./hybi/message'),
|
||||
Reader = require('./hybi/stream_reader');
|
||||
Message = require('./hybi/message');
|
||||
|
||||
var Hybi = function(request, url, options) {
|
||||
Base.apply(this, arguments);
|
||||
|
||||
this._extensions = new Extensions();
|
||||
this._reader = new Reader();
|
||||
this._stage = 0;
|
||||
this._masking = this._options.masking;
|
||||
this._protocols = this._options.protocols || [];
|
||||
@@ -62,14 +60,13 @@ Hybi.generateAccept = function(key) {
|
||||
Hybi.GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11';
|
||||
|
||||
var instance = {
|
||||
BYTE: 255,
|
||||
FIN: 128,
|
||||
MASK: 128,
|
||||
RSV1: 64,
|
||||
RSV2: 32,
|
||||
RSV3: 16,
|
||||
OPCODE: 15,
|
||||
LENGTH: 127,
|
||||
FIN: 0x80,
|
||||
MASK: 0x80,
|
||||
RSV1: 0x40,
|
||||
RSV2: 0x20,
|
||||
RSV3: 0x10,
|
||||
OPCODE: 0x0F,
|
||||
LENGTH: 0x7F,
|
||||
|
||||
OPCODES: {
|
||||
continuation: 0,
|
||||
@@ -84,8 +81,6 @@ var instance = {
|
||||
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) }),
|
||||
|
||||
ERRORS: {
|
||||
normal_closure: 1000,
|
||||
going_away: 1001,
|
||||
@@ -110,8 +105,8 @@ var instance = {
|
||||
return true;
|
||||
},
|
||||
|
||||
parse: function(data) {
|
||||
this._reader.put(data);
|
||||
parse: function(chunk) {
|
||||
this._reader.put(chunk);
|
||||
var buffer = true;
|
||||
while (buffer) {
|
||||
switch (this._stage) {
|
||||
@@ -133,16 +128,16 @@ var instance = {
|
||||
case 3:
|
||||
buffer = this._reader.read(4);
|
||||
if (buffer) {
|
||||
this._frame.maskingKey = buffer;
|
||||
this._stage = 4;
|
||||
this._frame.maskingKey = buffer;
|
||||
}
|
||||
break;
|
||||
|
||||
case 4:
|
||||
buffer = this._reader.read(this._frame.length);
|
||||
if (buffer) {
|
||||
this._emitFrame(buffer);
|
||||
this._stage = 0;
|
||||
this._emitFrame(buffer);
|
||||
}
|
||||
break;
|
||||
|
||||
@@ -169,6 +164,12 @@ var instance = {
|
||||
return this.frame(message, 'ping');
|
||||
},
|
||||
|
||||
pong: function(message) {
|
||||
if (this.readyState > 1) return false;
|
||||
message = message ||'';
|
||||
return this.frame(message, 'pong');
|
||||
},
|
||||
|
||||
close: function(reason, code) {
|
||||
reason = reason || '';
|
||||
code = code || this.ERRORS.normal_closure;
|
||||
@@ -186,27 +187,26 @@ var instance = {
|
||||
}
|
||||
},
|
||||
|
||||
frame: function(data, type, code) {
|
||||
if (this.readyState <= 0) return this._queue([data, type, code]);
|
||||
frame: function(buffer, type, code) {
|
||||
if (this.readyState <= 0) return this._queue([buffer, type, code]);
|
||||
if (this.readyState > 2) return false;
|
||||
|
||||
if (data instanceof Array) data = new Buffer(data);
|
||||
if (buffer instanceof Array) buffer = new Buffer(buffer);
|
||||
|
||||
var message = new Message(),
|
||||
isText = (typeof data === 'string'),
|
||||
payload, buffer;
|
||||
isText = (typeof buffer === 'string'),
|
||||
payload, copy;
|
||||
|
||||
message.rsv1 = message.rsv2 = message.rsv3 = false;
|
||||
message.opcode = this.OPCODES[type || (isText ? 'text' : 'binary')];
|
||||
|
||||
payload = isText ? new Buffer(data, 'utf8') : data;
|
||||
payload = isText ? new Buffer(buffer, 'utf8') : buffer;
|
||||
|
||||
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);
|
||||
copy = payload;
|
||||
payload = new Buffer(2 + copy.length);
|
||||
payload.writeUInt16BE(code, 0);
|
||||
copy.copy(payload, 2);
|
||||
}
|
||||
message.data = payload;
|
||||
|
||||
@@ -243,7 +243,6 @@ var instance = {
|
||||
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) |
|
||||
@@ -256,18 +255,11 @@ var instance = {
|
||||
buffer[1] = masked | length;
|
||||
} else if (length <= 65535) {
|
||||
buffer[1] = masked | 126;
|
||||
buffer[2] = ~~(length / 256);
|
||||
buffer[3] = length & BYTE;
|
||||
buffer.writeUInt16BE(length, 2);
|
||||
} 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;
|
||||
buffer.writeUInt32BE(Math.floor(length / 0x100000000), 2);
|
||||
buffer.writeUInt32BE(length % 0x100000000, 6);
|
||||
}
|
||||
|
||||
if (frame.masked) {
|
||||
@@ -295,39 +287,41 @@ var instance = {
|
||||
return new Buffer(headers.join('\r\n'), 'utf8');
|
||||
},
|
||||
|
||||
_shutdown: function(code, reason) {
|
||||
_shutdown: function(code, reason, error) {
|
||||
delete this._frame;
|
||||
delete this._message;
|
||||
this._stage = 5;
|
||||
|
||||
var sendCloseFrame = (this.readyState === 1);
|
||||
this.readyState = 2;
|
||||
this._stage = 5;
|
||||
|
||||
this._extensions.close(function() {
|
||||
if (sendCloseFrame) this.frame(reason, 'close', code);
|
||||
this.readyState = 3;
|
||||
if (error) this.emit('error', new Error(reason));
|
||||
this.emit('close', new Base.CloseEvent(code, reason));
|
||||
}, this);
|
||||
},
|
||||
|
||||
_fail: function(type, message) {
|
||||
if (this.readyState > 1) return;
|
||||
this.emit('error', new Error(message));
|
||||
this._shutdown(this.ERRORS[type], message);
|
||||
this._shutdown(this.ERRORS[type], message, true);
|
||||
},
|
||||
|
||||
_parseOpcode: function(data) {
|
||||
_parseOpcode: function(octet) {
|
||||
var rsvs = [this.RSV1, this.RSV2, this.RSV3].map(function(rsv) {
|
||||
return (data & rsv) === rsv;
|
||||
return (octet & rsv) === rsv;
|
||||
});
|
||||
|
||||
var frame = this._frame = new Frame();
|
||||
|
||||
frame.final = (data & this.FIN) === this.FIN;
|
||||
frame.final = (octet & this.FIN) === this.FIN;
|
||||
frame.rsv1 = rsvs[0];
|
||||
frame.rsv2 = rsvs[1];
|
||||
frame.rsv3 = rsvs[2];
|
||||
frame.opcode = (data & this.OPCODE);
|
||||
frame.opcode = (octet & this.OPCODE);
|
||||
|
||||
this._stage = 1;
|
||||
|
||||
if (!this._extensions.validFrameRsv(frame))
|
||||
return this._fail('protocol_error',
|
||||
@@ -343,38 +337,35 @@ var instance = {
|
||||
|
||||
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) {
|
||||
_parseLength: function(octet) {
|
||||
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');
|
||||
|
||||
frame.length = (data & this.LENGTH);
|
||||
frame.masked = (octet & this.MASK) === this.MASK;
|
||||
frame.length = (octet & this.LENGTH);
|
||||
|
||||
if (frame.length >= 0 && frame.length <= 125) {
|
||||
if (!this._checkFrameLength()) return;
|
||||
this._stage = frame.masked ? 3 : 4;
|
||||
if (!this._checkFrameLength()) return;
|
||||
} else {
|
||||
frame.lengthBytes = (frame.length === 126 ? 2 : 8);
|
||||
this._stage = 2;
|
||||
frame.lengthBytes = (frame.length === 126 ? 2 : 8);
|
||||
}
|
||||
|
||||
if (this._requireMasking && !frame.masked)
|
||||
return this._fail('unacceptable', 'Received unmasked frame but masking is required');
|
||||
},
|
||||
|
||||
_parseExtendedLength: function(buffer) {
|
||||
var frame = this._frame;
|
||||
frame.length = this._getInteger(buffer);
|
||||
frame.length = this._readUInt(buffer);
|
||||
|
||||
this._stage = frame.masked ? 3 : 4;
|
||||
|
||||
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 = frame.masked ? 3 : 4;
|
||||
},
|
||||
|
||||
_checkFrameLength: function() {
|
||||
@@ -467,11 +458,11 @@ var instance = {
|
||||
return buffer.toString('utf8', 0, buffer.length);
|
||||
},
|
||||
|
||||
_getInteger: function(bytes) {
|
||||
var number = 0;
|
||||
for (var i = 0, n = bytes.length; i < n; i++)
|
||||
number += bytes[i] * this.TWO_POWERS[n - 1 - i];
|
||||
return number;
|
||||
_readUInt: function(buffer) {
|
||||
if (buffer.length === 2) return buffer.readUInt16BE(0);
|
||||
|
||||
return buffer.readUInt32BE(0) * 0x100000000 +
|
||||
buffer.readUInt32BE(4);
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
var Stream = require('stream').Stream,
|
||||
url = require('url'),
|
||||
util = require('util'),
|
||||
Base = require('./base'),
|
||||
Headers = require('./headers'),
|
||||
HttpParser = require('../http_parser');
|
||||
|
||||
@@ -69,7 +70,7 @@ var instance = {
|
||||
this.headers = this._http.headers;
|
||||
|
||||
if (this.statusCode === 200) {
|
||||
this.emit('connect');
|
||||
this.emit('connect', new Base.ConnectEvent());
|
||||
} else {
|
||||
var message = "Can't establish a connection to the server at " + this._origin.href;
|
||||
this.emit('error', new Error(message));
|
||||
|
||||
@@ -21,10 +21,10 @@ var instance = {
|
||||
this.on('error', function() {});
|
||||
},
|
||||
|
||||
parse: function(data) {
|
||||
if (this._delegate) return this._delegate.parse(data);
|
||||
parse: function(chunk) {
|
||||
if (this._delegate) return this._delegate.parse(chunk);
|
||||
|
||||
this._http.parse(data);
|
||||
this._http.parse(chunk);
|
||||
if (!this._http.isComplete()) return;
|
||||
|
||||
this.method = this._http.method;
|
||||
|
||||
+27
-10
@@ -3,6 +3,7 @@
|
||||
var StreamReader = function() {
|
||||
this._queue = [];
|
||||
this._queueSize = 0;
|
||||
this._offset = 0;
|
||||
};
|
||||
|
||||
StreamReader.prototype.put = function(buffer) {
|
||||
@@ -15,13 +16,15 @@ StreamReader.prototype.put = function(buffer) {
|
||||
StreamReader.prototype.read = function(length) {
|
||||
if (length > this._queueSize) return null;
|
||||
if (length === 0) return new Buffer(0);
|
||||
|
||||
var queue = this._queue,
|
||||
first = queue[0],
|
||||
buffer;
|
||||
|
||||
|
||||
this._queueSize -= length;
|
||||
|
||||
var queue = this._queue,
|
||||
remain = length,
|
||||
first = queue[0],
|
||||
buffers, buffer;
|
||||
|
||||
if (first.length >= length) {
|
||||
this._queueSize -= length;
|
||||
if (first.length === length) {
|
||||
return queue.shift();
|
||||
} else {
|
||||
@@ -30,10 +33,8 @@ StreamReader.prototype.read = function(length) {
|
||||
return buffer;
|
||||
}
|
||||
}
|
||||
|
||||
var remain = length, buffers;
|
||||
|
||||
for (var i=0, n = queue.length; i < n; i++) {
|
||||
for (var i = 0, n = queue.length; i < n; i++) {
|
||||
if (remain < queue[i].length) break;
|
||||
remain -= queue[i].length;
|
||||
}
|
||||
@@ -43,10 +44,26 @@ StreamReader.prototype.read = function(length) {
|
||||
buffers.push(queue[0].slice(0, remain));
|
||||
queue[0] = queue[0].slice(remain);
|
||||
}
|
||||
this._queueSize -= length;
|
||||
return this._concat(buffers, length);
|
||||
};
|
||||
|
||||
StreamReader.prototype.eachByte = function(callback, context) {
|
||||
var buffer, n, index;
|
||||
|
||||
while (this._queue.length > 0) {
|
||||
buffer = this._queue[0];
|
||||
n = buffer.length;
|
||||
|
||||
while (this._offset < n) {
|
||||
index = this._offset;
|
||||
this._offset += 1;
|
||||
callback.call(context, buffer[index]);
|
||||
}
|
||||
this._offset = 0;
|
||||
this._queue.shift();
|
||||
}
|
||||
};
|
||||
|
||||
StreamReader.prototype._concat = function(buffers, length) {
|
||||
if (Buffer.concat) return Buffer.concat(buffers, length);
|
||||
|
||||
@@ -87,13 +87,13 @@ HttpParser.prototype.isComplete = function() {
|
||||
return this._complete;
|
||||
};
|
||||
|
||||
HttpParser.prototype.parse = function(data) {
|
||||
HttpParser.prototype.parse = function(chunk) {
|
||||
var offset = (version < 6) ? 1 : 0,
|
||||
consumed = this._parser.execute(data, 0, data.length) + offset;
|
||||
consumed = this._parser.execute(chunk, 0, chunk.length) + offset;
|
||||
|
||||
if (this._complete)
|
||||
this.body = (consumed < data.length)
|
||||
? data.slice(consumed)
|
||||
this.body = (consumed < chunk.length)
|
||||
? chunk.slice(consumed)
|
||||
: new Buffer(0);
|
||||
};
|
||||
|
||||
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
, "keywords" : ["websocket"]
|
||||
, "license" : "MIT"
|
||||
|
||||
, "version" : "0.5.4"
|
||||
, "version" : "0.6.1"
|
||||
, "engines" : {"node": ">=0.6.0"}
|
||||
, "main" : "./lib/websocket/driver"
|
||||
, "dependencies" : {"websocket-extensions": ">=0.1.1"}
|
||||
|
||||
@@ -184,7 +184,7 @@ test.describe("Client", function() { with(this) {
|
||||
}})
|
||||
|
||||
it("emits a 'connect' event when the proxy connects", function() { with(this) {
|
||||
expect(proxy, "emit").given("connect")
|
||||
expect(proxy, "emit").given("connect", anything())
|
||||
expect(proxy, "emit").given("close")
|
||||
expect(proxy, "emit").given("end")
|
||||
proxy.write(new Buffer("HTTP/1.1 200 OK\r\n\r\n"))
|
||||
|
||||
@@ -42,6 +42,28 @@ test.describe("draft-75", function() { with(this) {
|
||||
driver().parse([0x6c, 0x6f, 0xff])
|
||||
assertEqual( "Hello", message )
|
||||
}})
|
||||
|
||||
describe("when a message listener throws an error", function() { with(this) {
|
||||
before(function() { with(this) {
|
||||
this.messages = []
|
||||
|
||||
driver().on("message", function(msg) {
|
||||
messages.push(msg.data)
|
||||
throw new Error("an error")
|
||||
})
|
||||
}})
|
||||
|
||||
it("is not trapped by the parser", function() { with(this) {
|
||||
var buffer = [0x00, 0x48, 0x65, 0x6c, 0x6c, 0x6f, 0xff]
|
||||
assertThrows(Error, function() { driver().parse(buffer) })
|
||||
}})
|
||||
|
||||
it("parses text frames without dropping input", function() { with(this) {
|
||||
try { driver().parse([0x00, 0x48, 0x65, 0x6c, 0x6c, 0x6f, 0xff, 0x00, 0x57]) } catch (e) {}
|
||||
try { driver().parse([0x6f, 0x72, 0x6c, 0x64, 0xff]) } catch (e) {}
|
||||
assertEqual( ["Hello", "World"], messages )
|
||||
}})
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("frame", function() { with(this) {
|
||||
|
||||
@@ -191,6 +191,30 @@ test.describe("Hybi", function() { with(this) {
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("pong", function() { with(this) {
|
||||
it("does not write to the socket", function() { with(this) {
|
||||
expect(driver().io, "emit").exactly(0)
|
||||
driver().pong()
|
||||
}})
|
||||
|
||||
it("returns true", function() { with(this) {
|
||||
assertEqual( true, driver().pong() )
|
||||
}})
|
||||
|
||||
it("queues the pong until the handshake has been sent", function() { with(this) {
|
||||
expect(driver().io, "emit").given("data", buffer(
|
||||
"HTTP/1.1 101 Switching Protocols\r\n" +
|
||||
"Upgrade: websocket\r\n" +
|
||||
"Connection: Upgrade\r\n" +
|
||||
"Sec-WebSocket-Accept: JdiiuafpBKRqD7eol0y4vJDTsTs=\r\n" +
|
||||
"\r\n"))
|
||||
expect(driver().io, "emit").given("data", buffer([0x8a, 0]))
|
||||
|
||||
driver().pong()
|
||||
driver().start()
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("close", function() { with(this) {
|
||||
it("does not write anything to the socket", function() { with(this) {
|
||||
expect(driver().io, "emit").exactly(0)
|
||||
@@ -361,6 +385,28 @@ test.describe("Hybi", function() { with(this) {
|
||||
driver().parse([0x89, 0x04, 0x4f, 0x48, 0x41, 0x49])
|
||||
assertEqual( [0x8a, 0x04, 0x4f, 0x48, 0x41, 0x49], collector().bytes )
|
||||
}})
|
||||
|
||||
describe("when a message listener throws an error", function() { with(this) {
|
||||
before(function() { with(this) {
|
||||
this.messages = []
|
||||
|
||||
driver().on("message", function(msg) {
|
||||
messages.push(msg.data)
|
||||
throw new Error("an error")
|
||||
})
|
||||
}})
|
||||
|
||||
it("is not trapped by the parser", function() { with(this) {
|
||||
var buffer = [0x81, 0x05, 0x48, 0x65, 0x6c, 0x6c, 0x6f]
|
||||
assertThrows(Error, function() { driver().parse(buffer) })
|
||||
}})
|
||||
|
||||
it("parses unmasked text frames without dropping input", function() { with(this) {
|
||||
try { driver().parse([0x81, 0x05, 0x48, 0x65, 0x6c, 0x6c, 0x6f, 0x81, 0x05]) } catch (e) {}
|
||||
try { driver().parse([0x57, 0x6f, 0x72, 0x6c, 0x64]) } catch (e) {}
|
||||
assertEqual( ["Hello", "World"], messages )
|
||||
}})
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("frame", function() { with(this) {
|
||||
@@ -425,6 +471,17 @@ test.describe("Hybi", function() { with(this) {
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("pong", function() { with(this) {
|
||||
it("writes a pong frame to the socket", function() { with(this) {
|
||||
driver().pong("mic check")
|
||||
assertEqual([0x8a, 0x09, 0x6d, 0x69, 0x63, 0x20, 0x63, 0x68, 0x65, 0x63, 0x6b], collector().bytes)
|
||||
}})
|
||||
|
||||
it("returns true", function() { with(this) {
|
||||
assertEqual(true, driver().pong())
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("close", function() { with(this) {
|
||||
it("writes a close frame to the socket", function() { with(this) {
|
||||
driver().close("<%= reasons %>", 1003)
|
||||
@@ -498,6 +555,17 @@ test.describe("Hybi", function() { with(this) {
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("pong", function() { with(this) {
|
||||
it("does not write to the socket", function() { with(this) {
|
||||
expect(driver().io, "emit").exactly(0)
|
||||
driver().pong()
|
||||
}})
|
||||
|
||||
it("returns false", function() { with(this) {
|
||||
assertEqual( false, driver().pong() )
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("close", function() { with(this) {
|
||||
it("does not write to the socket", function() { with(this) {
|
||||
expect(driver().io, "emit").exactly(0)
|
||||
@@ -558,6 +626,17 @@ test.describe("Hybi", function() { with(this) {
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("pong", function() { with(this) {
|
||||
it("does not write to the socket", function() { with(this) {
|
||||
expect(driver().io, "emit").exactly(0)
|
||||
driver().pong()
|
||||
}})
|
||||
|
||||
it("returns false", function() { with(this) {
|
||||
assertEqual( false, driver().pong() )
|
||||
}})
|
||||
}})
|
||||
|
||||
describe("close", function() { with(this) {
|
||||
it("does not write to the socket", function() { with(this) {
|
||||
expect(driver().io, "emit").exactly(0)
|
||||
|
||||
Reference in New Issue
Block a user