refactoring message broker into plugin
This commit is contained in:
parent
8ae205de6c
commit
8d672de5bd
@ -3,9 +3,6 @@
|
|||||||
// server-side socket behaviour
|
// server-side socket behaviour
|
||||||
var ios = null; // io is already taken in express
|
var ios = null; // io is already taken in express
|
||||||
var util = require('bitcore').util;
|
var util = require('bitcore').util;
|
||||||
var mdb = require('../../lib/MessageDb').default();
|
|
||||||
var microtime = require('microtime');
|
|
||||||
var enableMessageBroker;
|
|
||||||
|
|
||||||
var verbose = false;
|
var verbose = false;
|
||||||
var log = function() {
|
var log = function() {
|
||||||
@ -15,64 +12,21 @@ var log = function() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
module.exports.init = function(io_ext, config) {
|
module.exports.init = function(io_ext, config) {
|
||||||
enableMessageBroker = config ? config.enableMessageBroker : false;
|
|
||||||
ios = io_ext;
|
ios = io_ext;
|
||||||
|
verbose = config.verbose;
|
||||||
if (ios) {
|
if (ios) {
|
||||||
// when a new socket connects
|
// when a new socket connects
|
||||||
ios.sockets.on('connection', function(socket) {
|
ios.sockets.on('connection', function(socket) {
|
||||||
log('New connection from ' + socket.id);
|
log('New connection from ' + socket.id);
|
||||||
// when it subscribes, make it join the according room
|
// when it subscribes, make it join the according room
|
||||||
socket.on('subscribe', function(topic) {
|
socket.on('subscribe', function(topic) {
|
||||||
if (socket.rooms.length === 1) {
|
log('subscribe to ' + topic);
|
||||||
log('subscribe to ' + topic);
|
socket.join(topic);
|
||||||
socket.join(topic);
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
|
|
||||||
if (enableMessageBroker) {
|
|
||||||
// when it requests sync, send him all pending messages
|
|
||||||
socket.on('sync', function(ts) {
|
|
||||||
log('Sync requested by ' + socket.id);
|
|
||||||
log(' from timestamp '+ts);
|
|
||||||
var rooms = socket.rooms;
|
|
||||||
if (rooms.length !== 2) {
|
|
||||||
socket.emit('insight-error', 'Must subscribe with public key before syncing');
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
var to = rooms[1];
|
|
||||||
var upper_ts = Math.round(microtime.now());
|
|
||||||
log(' to timestamp '+upper_ts);
|
|
||||||
mdb.getMessages(to, ts, upper_ts, function(err, messages) {
|
|
||||||
if (err) {
|
|
||||||
throw new Error('Couldn\'t get messages on sync request: ' + err);
|
|
||||||
}
|
|
||||||
log('\tFound ' + messages.length + ' message' + (messages.length !== 1 ? 's' : ''));
|
|
||||||
for (var i = 0; i < messages.length; i++) {
|
|
||||||
broadcastMessage(messages[i], socket);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
// when it sends a message, add it to db
|
|
||||||
socket.on('message', function(m) {
|
|
||||||
log('Message sent from ' + m.pubkey + ' to ' + m.to);
|
|
||||||
mdb.addMessage(m, function(err) {
|
|
||||||
if (err) {
|
|
||||||
throw new Error('Couldn\'t add message to database: ' + err);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
|
|
||||||
// disconnect handler
|
|
||||||
socket.on('disconnect', function() {
|
|
||||||
log('disconnected ' + socket.id);
|
|
||||||
});
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
if (enableMessageBroker)
|
|
||||||
mdb.on('message', broadcastMessage);
|
|
||||||
}
|
}
|
||||||
|
return ios;
|
||||||
};
|
};
|
||||||
|
|
||||||
var simpleTx = function(tx) {
|
var simpleTx = function(tx) {
|
||||||
|
|||||||
@ -61,6 +61,7 @@ dataDir += network === 'testnet' ? 'testnet3' : '';
|
|||||||
|
|
||||||
var safeConfirmations = process.env.INSIGHT_SAFE_CONFIRMATIONS || 6;
|
var safeConfirmations = process.env.INSIGHT_SAFE_CONFIRMATIONS || 6;
|
||||||
var ignoreCache = process.env.INSIGHT_IGNORE_CACHE || 0;
|
var ignoreCache = process.env.INSIGHT_IGNORE_CACHE || 0;
|
||||||
|
var verbose = process.env.VERBOSE === 'true';
|
||||||
|
|
||||||
|
|
||||||
var bitcoindConf = {
|
var bitcoindConf = {
|
||||||
@ -76,7 +77,7 @@ var bitcoindConf = {
|
|||||||
disableAgent: true
|
disableAgent: true
|
||||||
};
|
};
|
||||||
|
|
||||||
var enableMessageBroker = process.env.ENABLE_MESSAGE_BROKER === 'true';
|
var enableMailbox = process.env.ENABLE_MAILBOX === 'true';
|
||||||
|
|
||||||
if (!fs.existsSync(db)) {
|
if (!fs.existsSync(db)) {
|
||||||
var err = fs.mkdirSync(db);
|
var err = fs.mkdirSync(db);
|
||||||
@ -89,7 +90,8 @@ if (!fs.existsSync(db)) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
enableMessageBroker: enableMessageBroker,
|
enableMailbox: enableMailbox,
|
||||||
|
verbose: verbose,
|
||||||
version: version,
|
version: version,
|
||||||
root: rootPath,
|
root: rootPath,
|
||||||
publicPath: process.env.INSIGHT_PUBLIC_PATH || false,
|
publicPath: process.env.INSIGHT_PUBLIC_PATH || false,
|
||||||
|
|||||||
53
insight.js
53
insight.js
@ -111,16 +111,63 @@ if (!config.disableHistoricSync) {
|
|||||||
if (peerSync) peerSync.allowReorgs = true;
|
if (peerSync) peerSync.allowReorgs = true;
|
||||||
|
|
||||||
|
|
||||||
//express settings
|
// express settings
|
||||||
require('./config/express')(expressApp, historicSync, peerSync);
|
require('./config/express')(expressApp, historicSync, peerSync);
|
||||||
|
|
||||||
//Bootstrap routes
|
// routes
|
||||||
require('./config/routes')(expressApp);
|
require('./config/routes')(expressApp);
|
||||||
|
|
||||||
// socket.io
|
// socket.io
|
||||||
var server = require('http').createServer(expressApp);
|
var server = require('http').createServer(expressApp);
|
||||||
var ios = require('socket.io')(server);
|
var ios = require('socket.io')(server);
|
||||||
require('./app/controllers/socket.js').init(ios, config);
|
|
||||||
|
// plugins
|
||||||
|
require('./app/controllers/socket.js').init(ios);
|
||||||
|
if (config.enableMessageBroker) {
|
||||||
|
var mdb = require('../../lib/MessageDb').default();
|
||||||
|
|
||||||
|
// when it requests sync, send him all pending messages
|
||||||
|
socket.on('sync', function(ts) {
|
||||||
|
log('Sync requested by ' + socket.id);
|
||||||
|
log(' from timestamp ' + ts);
|
||||||
|
var rooms = socket.rooms;
|
||||||
|
if (rooms.length !== 2) {
|
||||||
|
socket.emit('insight-error', 'Must subscribe with public key before syncing');
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
var to = rooms[1];
|
||||||
|
var upper_ts = Math.round(microtime.now());
|
||||||
|
log(' to timestamp ' + upper_ts);
|
||||||
|
mdb.getMessages(to, ts, upper_ts, function(err, messages) {
|
||||||
|
if (err) {
|
||||||
|
throw new Error('Couldn\'t get messages on sync request: ' + err);
|
||||||
|
}
|
||||||
|
log('\tFound ' + messages.length + ' message' + (messages.length !== 1 ? 's' : ''));
|
||||||
|
for (var i = 0; i < messages.length; i++) {
|
||||||
|
broadcastMessage(messages[i], socket);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// when it sends a message, add it to db
|
||||||
|
socket.on('message', function(m) {
|
||||||
|
log('Message sent from ' + m.pubkey + ' to ' + m.to);
|
||||||
|
mdb.addMessage(m, function(err) {
|
||||||
|
if (err) {
|
||||||
|
throw new Error('Couldn\'t add message to database: ' + err);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
|
||||||
|
// disconnect handler
|
||||||
|
socket.on('disconnect', function() {
|
||||||
|
log('disconnected ' + socket.id);
|
||||||
|
});
|
||||||
|
|
||||||
|
mdb.on('message', broadcastMessage);
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
//Start the app by listening on <port>
|
//Start the app by listening on <port>
|
||||||
server.listen(config.port, function() {
|
server.listen(config.port, function() {
|
||||||
|
|||||||
2
plugins/mailbox.js
Normal file
2
plugins/mailbox.js
Normal file
@ -0,0 +1,2 @@
|
|||||||
|
|
||||||
|
var microtime = require('microtime');
|
||||||
@ -15,7 +15,7 @@ describe('socket server', function() {
|
|||||||
});
|
});
|
||||||
it('should register socket handlers', function() {
|
it('should register socket handlers', function() {
|
||||||
var io = {
|
var io = {
|
||||||
sockets: new EventEmitter()
|
sockets: new EventEmitter(),
|
||||||
}
|
}
|
||||||
socket.init(io);
|
socket.init(io);
|
||||||
|
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user