servicebus
Advanced tools
Comparing version 0.3.1 to 0.3.2
@@ -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 () {}; |
33639
4
744
+ Addedbluebird@~1.0.0
+ Addedbluebird@1.0.8(transitive)