diff --git a/lib/client.js b/lib/client.js index d621b76..1cef2fc 100644 --- a/lib/client.js +++ b/lib/client.js @@ -2,7 +2,7 @@ var Stream = require('stream').Stream; var EventEmitter = require('events').EventEmitter; var BufferReadStream = require('streamers').BufferReadStream; -var WebSocket = require('streamws'); +var WebSocket = require('ws'); var util = require('./util'); var BinaryStream = require('./stream').BinaryStream; @@ -24,18 +24,24 @@ function BinaryClient(socket, options) { if(typeof socket === 'string') { this._nextId = 0; - this._socket = new WebSocket(socket); + this._socket = new WebSocket(socket, this._options.ws || {}); } else { // Use odd numbered ids for server originated streams this._nextId = 1; this._socket = socket; } - this._socket.binaryType = 'arraybuffer'; + //this._socket.binaryType = 'arraybuffer'; this._socket.addEventListener('open', function(){ self.emit('open'); }); + this._socket.addEventListener('ping', function(data, flags){ + self.emit('ping', data, flags); + }); + this._socket.addEventListener('pong', function(data, flags){ + self.emit('pong', data, flags); + }); // if node this._socket.on('drain', function(){ var ids = Object.keys(self.streams); @@ -59,7 +65,8 @@ function BinaryClient(socket, options) { self.emit('close', code, message); }); this._socket.addEventListener('message', function(data, flags){ - util.setZeroTimeout(function(){ + //self.lastActive = Date.now(); + //util.setZeroTimeout(function(){ // Message format // [type, payload, bonus ] @@ -90,7 +97,13 @@ function BinaryClient(socket, options) { // // Close // [ 6 , null , streamId ] - // + // + // 7 + // [ 7 , Notification , streamId ] + // + // 8 + // [ 8 , Message ] + // data = data.data; @@ -115,6 +128,9 @@ function BinaryClient(socket, options) { var streamId = data[2]; var binaryStream = self._receiveStream(streamId); self.emit('stream', binaryStream, meta); + if(meta.head) { + binaryStream._onData(meta.head); + } break; case 2: var payload = data[1]; @@ -123,7 +139,7 @@ function BinaryClient(socket, options) { if(binaryStream) { binaryStream._onData(payload); } else { - self.emit('error', new Error('Received `data` message for unknown stream: ' + streamId)); + self.emit('warn', new Error('Received `data` message for unknown stream: ' + streamId)); } break; case 3: @@ -132,7 +148,7 @@ function BinaryClient(socket, options) { if(binaryStream) { binaryStream._onPause(); } else { - self.emit('error', new Error('Received `pause` message for unknown stream: ' + streamId)); + self.emit('warn', new Error('Received `pause` message for unknown stream: ' + streamId)); } break; case 4: @@ -141,7 +157,7 @@ function BinaryClient(socket, options) { if(binaryStream) { binaryStream._onResume(); } else { - self.emit('error', new Error('Received `resume` message for unknown stream: ' + streamId)); + self.emit('warn', new Error('Received `resume` message for unknown stream: ' + streamId)); } break; case 5: @@ -150,7 +166,7 @@ function BinaryClient(socket, options) { if(binaryStream) { binaryStream._onEnd(); } else { - self.emit('error', new Error('Received `end` message for unknown stream: ' + streamId)); + self.emit('warn', new Error('Received `end` message for unknown stream: ' + streamId)); } break; case 6: @@ -159,18 +175,45 @@ function BinaryClient(socket, options) { if(binaryStream) { binaryStream._onClose(); } else { - self.emit('error', new Error('Received `close` message for unknown stream: ' + streamId)); + self.emit('warn', new Error('Received `close` message for unknown stream: ' + streamId)); + } + break; + case 7: + var streamId = data[2]; + var binaryStream = self.streams[streamId]; + if(binaryStream) { + var event = data[1]; + binaryStream._onNotification(event); + } else { + self.emit('warn', new Error('Received `notification` message for unknown stream: ' + streamId)); } break; + case 8: + var message = data[1]; + self.emit('message', message); + break; default: - self.emit('error', new Error('Unrecognized message type received: ' + data[0])); + self.emit('warn', new Error('Unrecognized message type received: ' + data[0])); } - }); + //}); }); } util.inherits(BinaryClient, EventEmitter); +BinaryClient.prototype.ping = function(data){ + return this._socket.ping(data, {}, true); +} + +BinaryClient.prototype.pong = function(data){ + return this._socket.pong(data, {}, true); +} + +BinaryClient.prototype.sendMessage = function(message){ + var data = util.pack([8, message, 0]); + return this._socket.send(data, {binary: true}); +} + BinaryClient.prototype.send = function(data, meta){ var stream = this.createStream(meta); if(data instanceof Stream) { diff --git a/lib/server.js b/lib/server.js index f02d1c4..a4032b3 100644 --- a/lib/server.js +++ b/lib/server.js @@ -1,4 +1,4 @@ -var ws = require('streamws'); +var ws = require('ws'); var EventEmitter = require('events').EventEmitter; var util = require('./util'); @@ -24,7 +24,7 @@ function BinaryServer(options) { this._server = new ws.Server(options); } - this._server.on('connection', function(socket){ + this._server.on('connection', function(socket, req){ var clientId = self._clientCounter; var binaryClient = new BinaryClient(socket, options); binaryClient.id = clientId; @@ -33,7 +33,7 @@ function BinaryServer(options) { binaryClient.on('close', function(){ delete self.clients[clientId]; }); - self.emit('connection', binaryClient); + self.emit('connection', binaryClient, req); }); this._server.on('error', function(error){ self.emit('error', error); diff --git a/lib/stream.js b/lib/stream.js index c0c3433..34239ac 100644 --- a/lib/stream.js +++ b/lib/stream.js @@ -20,6 +20,7 @@ function BinaryStream(socket, id, create, meta) { this._closed = false; this._ended = false; + this._buf = []; if(create) { // This is a stream we are creating @@ -61,11 +62,26 @@ BinaryStream.prototype._onPause = function() { this.emit('pause'); }; +BinaryStream.prototype.flow = function() { + var data = this._buf.shift(); + if (data) { + this.emit('data', data); + } + if (!this.paused && this._buf.length > 0) { + process.nextTick(() => { + this.flow(); + }) + } +}; + BinaryStream.prototype._onResume = function() { // Emit resume event this.paused = false; this.emit('resume'); this.emit('drain'); + process.nextTick(()=>{ + this.flow(); + }) }; BinaryStream.prototype._write = function(code, data, bonus) { @@ -73,7 +89,7 @@ BinaryStream.prototype._write = function(code, data, bonus) { return false; } var message = util.pack([code, data, bonus]); - return this._socket.send(message) !== false; + return this._socket.send(message, {binary: true}) !== false; }; BinaryStream.prototype.write = function(data) { @@ -110,8 +126,19 @@ BinaryStream.prototype._onEnd = function() { }; BinaryStream.prototype._onData = function(data) { - // Dispatch - this.emit('data', data); + if(!this.paused && this._buf.length === 0 && this.listenerCount('data') > 0) { + this.emit('data', data); + } else { + this._buf.push(data); + } +}; + +BinaryStream.prototype._onNotification = function(event) { + this.emit('notify', event); +}; + +BinaryStream.prototype.sendNotification = function(event) { + this._write(7, event, this.id); }; BinaryStream.prototype.pause = function() { diff --git a/package.json b/package.json index 8546c78..ff47c20 100644 --- a/package.json +++ b/package.json @@ -2,12 +2,8 @@ "author": "Eric Zhang (http://ericzhang.com)", "name": "binaryjs", "description": "Binary realtime streaming made easy", - "version": "0.2.2", + "version": "0.2.3", "homepage": "http://binaryjs.com", - "repository": { - "type": "git", - "url": "git://github.com/binaryjs/binaryjs.git" - }, "main": "lib/server.js", "scripts": { "test": "make test" @@ -16,12 +12,12 @@ "node": ">=0.10.20" }, "dependencies": { - "streamws": ">=0.1.1", - "binarypack": ">=0.0.4", + "ws": "https://github.com/carck/ws#9ada0bbf7e4a6e49744d8457c48ac253b3f7a5c6", + "binarypack": "https://github.com/carck/node-binarypack#0eed17aa94ef5d9fd9ec67e983a12631f57bfa0b", "streamers": ">=0.1.0" }, "devDependencies": { "mocha": "~1.3.0", "uglify-js": "~1.3.5" } -} \ No newline at end of file +} diff --git a/test/client.js b/test/client.js index 2d13b24..449336b 100644 --- a/test/client.js +++ b/test/client.js @@ -37,6 +37,21 @@ describe('BinaryClient', function(){ client.createStream(); }); }); + it('should receive notification', function(done){ + server.on('connection', function(client){ + client.on('stream', function(stream){ + stream.sendNotification("ready"); + }); + }); + var client = new BinaryClient(serverUrl); + client.on('open', function(){ + client.createStream().on('notify', function(event){ + if('ready' === event) { + done(); + } + }); + }); + }); }); describe('sending data', function(){