Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion index.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@ module.exports = {
createConsumer: consumerService.createConsumer,
createMessage: messageService.createMessage,
createTopic: topicService.createTopic,
getConsumersForTopic: topicService.getConsumersForTopic,
getTopics: topicService.getTopics,
getStatus: queueService.getStatus,
getProcessedMessages: queueService.getProcessedMessages
getProcessedMessages: queueService.getProcessedMessages,
getQueueMessages: queueService.getQueueMessages
};
9 changes: 9 additions & 0 deletions models/message.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ class Message {
this._value = '';
this._processing_details = [];
this._processed = false;
this._dropped = false;
this._allowed_retries = 0;
}

Expand Down Expand Up @@ -49,6 +50,14 @@ class Message {
this._processed = value;
}

getDropped() {
return this._dropped;
}

setDropped(value) {
this._dropped = value;
}

setAllowedRetries(value) {
this._allowed_retries = value;
}
Expand Down
6 changes: 6 additions & 0 deletions models/queue.js
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ class Queue {
this._retries = retries;
this._queue = [];
this._processed_messages = [];
this._messages = [];
}

getSize() {
Expand All @@ -27,6 +28,7 @@ class Queue {

enQueue(message) {
this._queue.unshift(message);
this._messages.push(message);
}

numberOfMessagesInQueue() {
Expand Down Expand Up @@ -57,6 +59,10 @@ class Queue {
setProcessedMessages(processedMessages) {
this._processed_messages = processedMessages;
}

getMessages() {
return this._messages;
}
}

module.exports = Queue;
1 change: 1 addition & 0 deletions services/message.handler.service.js
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ const processMessage = function (message) {
if (_.isEmpty(consumersForTopic)) {
logger.info(`No consumers found for topic ${message.getTopic()}`);
message.setProcessed(true);
message.setDropped(true);
return Promise.resolve(message);
} else {
const consumersByPriority = _.groupBy(consumersForTopic, consumer => {
Expand Down
15 changes: 14 additions & 1 deletion services/queue.service.js
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,7 @@ const getProcessedMessages = function () {
}
return Promise.resolve(QueueInstance.getProcessedMessages());
};

const pushMessageToProcessedMessage = function (message) {
if (!QueueInstance) {
try {
Expand All @@ -106,9 +107,21 @@ const pushMessageToProcessedMessage = function (message) {
return Promise.resolve();
};

const getQueueMessages = function() {
if (!QueueInstance) {
try {
initQueue();
} catch (err) {
return Promise.reject(err.toString());
}
}
return Promise.resolve(QueueInstance.getMessages());
};

module.exports = {
pushMessageToQueue,
getStatus,
getProcessedMessages,
pushMessageToProcessedMessage
pushMessageToProcessedMessage,
getQueueMessages,
};
11 changes: 10 additions & 1 deletion services/topic.service.js
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,19 @@ const getConsumersForTopic = function (topic) {
return Promise.resolve(topicVsConsumer[topic]);
};

const getTopics = function () {
const topics = Object.keys(topicVsConsumer);
if (topics.length === 0) {
return Promise.reject(`No topics created`);
} else {
return Promise.resolve(topics);
}
}

module.exports = {
isValidTopic,
createTopic,
registerConsumerForTopic,
getConsumersForTopic
getConsumersForTopic,
getTopics
};