Huge News!Announcing our $40M Series B led by Abstract Ventures.Learn More
Socket
Sign inDemoInstall
Socket

servicebus

Package Overview
Dependencies
Maintainers
2
Versions
73
Alerts
File Explorer

Advanced tools

Socket logo

Install Socket

Detect and block malicious and high-risk dependencies

Install

servicebus - npm Package Compare versions

Comparing version 0.3.1 to 0.3.2

makefile

14

bus/bus.js

@@ -13,3 +13,4 @@ function Bus () {

Bus.prototype.handleIncoming = function (message, headers, deliveryInfo, messageHandle, options, callback) {
var stack = this.incomingMiddleware, index = this.incomingMiddleware.length - 1;
var index = this.incomingMiddleware.length - 1;
var self = this;

@@ -27,3 +28,3 @@ function next (err) {

layer = stack[index--];
layer = self.incomingMiddleware[index--];

@@ -33,3 +34,3 @@ if ( undefined === layer) {

} else {
layer(message, headers, deliveryInfo, messageHandle, options, next);
layer.call(self, message, headers, deliveryInfo, messageHandle, options, next);
}

@@ -43,3 +44,4 @@ }

var stack = this.outgoingMiddleware, index = 0;
var index = 0;
var self = this;

@@ -54,3 +56,3 @@ function next (err) {

layer = stack[index];
layer = self.outgoingMiddleware[index];

@@ -62,3 +64,3 @@ index++;

} else {
layer(queueName, message, next);
layer.call(self, queueName, message, next);
}

@@ -65,0 +67,0 @@ }

@@ -1,26 +0,5 @@

var rejected = {}, maxRetries = 10;
var log = require('debug')('servicebus:retry');
var maxRetries = 10, rejected = {};
var util = require('util');
function retry (message, headers, deliveryInfo, messageHandle, options, next) {
if (options && options.ack) {
message.handle = {
ack: function () { messageHandle.acknowledge(); },
acknowledge: function () { messageHandle.acknowledge(); },
reject: function () {
var msgRejected = rejected[message.cid] || 0;
if (msgRejected >= maxRetries) {
messageHandle.acknowledge();
delete rejected[message.cid];
} else {
msgRejected++;
rejected[message.cid] = msgRejected;
messageHandle.reject(true);
}
}
};
}
next(null, message, headers, deliveryInfo, messageHandle, options);
}
module.exports = function (options) {

@@ -32,4 +11,36 @@ options = options || {};

return {
handleIncoming: retry
handleIncoming: function retry (message, headers, deliveryInfo, messageHandle, options, next) {
var bus = this;
if (options && options.ack) {
message.handle = {
ack: function () { messageHandle.acknowledge(); },
acknowledge: function () { messageHandle.acknowledge(); },
reject: function () {
if (rejected[message.cid] == undefined) {
rejected[message.cid] = 1
} else {
rejected[message.cid] = rejected[message.cid] + 1;
}
if (rejected[message.cid] > maxRetries) {
var errorQueueName = util.format('%s.error', deliveryInfo.queue);
log('sending message %s to error queue %s', message.cid, errorQueueName);
bus.connection.publish(errorQueueName, message, { contentType: 'application/json', deliveryMode: 2 });
messageHandle.acknowledge();
delete rejected[message.cid];
} else {
log('retrying message %s', message.cid);
messageHandle.reject(true);
}
}
};
}
next(null, message, headers, deliveryInfo, messageHandle, options);
}
};
}

@@ -7,2 +7,3 @@ var amqp = require('amqp'),

PubSubQueue = require('./pubsubqueue'),
Promise = require('bluebird'),
Queue = require('./queue'),

@@ -19,3 +20,2 @@ util = require('util');

this.delayOnStartup = options.delayOnStartup || 10;
this.initialized = false;
this.log = options.log || log;

@@ -31,7 +31,2 @@ this.pubsubqueues = {};

this.connection.on('error', function (err) {
self.log('Error connecting to rabbitmq at ' + options.url + ' error: ' + err.toString());
throw err;
});
this.connection.on('close', function () {

@@ -41,5 +36,14 @@ self.log('rabbitmq connection closed.');

this.connection.on('ready', function () {
self.initialized = true;
self.log("rabbitmq connected to " + self.connection.serverProperties.product);
this.initialized = new Promise(function (resolve, reject) {
self.connection.on('error', reject);
self.connection.on('ready', function () {
self.log("rabbitmq connected to " + self.connection.serverProperties.product);
resolve();
});
}).catch(Promise.CancellationError, function (err) {
self.log('Error connecting to rabbitmq at ' + options.url + ' error: ' + err.toString());
throw err;
});

@@ -63,3 +67,3 @@

if (self.initialized) {
this.initialized.done(function() {

@@ -71,11 +75,4 @@ if (self.queues[queueName] === undefined) {

self.queues[queueName].listen(callback, options);
} else {
self.connection.on('ready', function() {
self.log(queueName + ' ready');
process.nextTick(function() {
self.initialized = true;
self.listen(queueName, options, callback);
});
});
}
});
};

@@ -86,3 +83,3 @@

if (self.initialized) {
this.initialized.done(function() {

@@ -96,18 +93,4 @@ if (self.queues[queueName] === undefined) {

});
} else {
var resend = function() {
self.initialized = true;
self.send(queueName, message);
};
var timeout = function(){
self.log('timout triggered');
self.connection.removeListener('ready', resend);
process.nextTick(resend);
};
var timeoutId = setTimeout(timeout, self.delayOnStartup);
self.connection.on('ready', function() {
clearTimeout(timeoutId);
process.nextTick(resend);
});
}
});
};

@@ -122,17 +105,12 @@

}
if (self.initialized) {
this.initialized.done(function() {
if (self.pubsubqueues[queueName] === undefined) {
if (self.pubsubqueues[queueName] === undefined) {
self.pubsubqueues[queueName] = new PubSubQueue({ bus: self, connection: self.connection, queueName: queueName, log: self.log });
}
self.pubsubqueues[queueName].subscribe(callback, options);
} else {
self.connection.on('ready', function() {
process.nextTick(function() {
self.initialized = true;
self.subscribe(queueName, options, callback);
});
});
}
self.pubsubqueues[queueName].subscribe(callback, options);
});
};

@@ -143,3 +121,3 @@

if (self.initialized) {
this.initialized.done(function() {

@@ -154,13 +132,7 @@ if (self.pubsubqueues[queueName] === undefined) {

});
} else {
var republish = function() {
self.initialized = true;
self.publish(queueName, message);
};
self.connection.on('ready', function() {
process.nextTick(republish);
});
}
});
};
module.exports.Bus = RabbitMQBus;

@@ -0,0 +0,0 @@ var Correlator = require('./correlator'),

@@ -5,7 +5,8 @@ {

"mattwalters5@gmail.com",
"timisbusy@gmail.com"
"timisbusy@gmail.com",
"mail@chrisabrams.com"
],
"name": "servicebus",
"description": "Simple service bus for sending events between processes using amqp.",
"version": "0.3.1",
"version": "0.3.2",
"homepage": "https://github.com/mateodelnorte/servicebus",

@@ -22,3 +23,4 @@ "repository": {

"node-uuid": "1.4.0",
"debug": "~0.7.2"
"debug": "~0.7.2",
"bluebird": "~1.0.0"
},

@@ -25,0 +27,0 @@ "devDependencies": {

@@ -32,3 +32,3 @@ # servicebus

(Note: message acking requires use of the retry() middleware, referenced below below)
(Note: message acking requires use of the retry() middleware, referenced below)

@@ -35,0 +35,0 @@ Servicebus integrates with RabbitMQ's message acknowledement functionality, which causes messages to queue instead of sending until the listening processes marks any previously received message as acknowledged or rejected. Messages can be acknowledged or rejected with the following syntax. To use ack and reject, it must be specified when defining the listening function:

require('longjohn')
if ( ! process.env.RABBITMQ_URL)
throw new Error('Tests require a RABBITMQ_URL environment variable to be set, pointing to the RabbiqMQ instance you wish to use.');
throw new Error('Tests require a RABBITMQ_URL environment variable to be set, pointing to the RabbiqMQ instance you wish to use. Example url: "amqp://localhost:5672"');

@@ -6,0 +6,0 @@ var busUrl = process.env.RABBITMQ_URL

@@ -0,0 +0,0 @@ var noop = function () {};

SocketSocket SOC 2 Logo

Product

  • Package Alerts
  • Integrations
  • Docs
  • Pricing
  • FAQ
  • Roadmap
  • Changelog

Packages

npm

Stay in touch

Get open source security insights delivered straight into your inbox.


  • Terms
  • Privacy
  • Security

Made with ⚡️ by Socket Inc