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
8 changes: 4 additions & 4 deletions examples/app.js
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ const httpPort = 8080;
const app = express();

// Using queue middleware
const queueMw = expressQueue({ activeLimit: 2, queuedLimit: 6 });
const queueMw = expressQueue({ activeLimit: 2, queuedLimit: 6, field:"userId" });
app.use(queueMw);
// May be also:
// app.use(queue({ activeLimit: 2, queuedLimit: -1 }));
Expand All @@ -26,16 +26,16 @@ let counter = 0;

app.get('/test1', function (req, res) {
let cnt = counter++; // local var inside the closure
console.log(`get(test1): [${cnt}/request] queueLength: ${queueMw.queue.getLength()}`);
console.log(`get(test1) ${req.query.userId}: [${cnt}/request] queueLength: ${queueMw.queue.getLength()}`);

const result = { test: 'test' };

setTimeout(function() {
console.log(`get(test1): [${cnt}/ready] queueLength: ${queueMw.queue.getLength()}` );
console.log(`get(test1) ${req.query.userId}: [${cnt}/ready] queueLength: ${queueMw.queue.getLength()}` );
res
.status(200)
.send(result);
console.log(`get(test1): [${cnt}/sent] queueLength: ${queueMw.queue.getLength()}`);
console.log(`get(test1) ${req.query.userId}: [${cnt}/sent] queueLength: ${queueMw.queue.getLength()}`);
}, RESPONSE_DELAY);

});
Expand Down
107 changes: 71 additions & 36 deletions lib/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -6,63 +6,98 @@ const endMw = require('express-end');
const MiniQueue = require('mini-queue');


const expressQueueMw = function(config) {
const expressQueueMw = function (config) {

debug('Initializing: config:', config);

const self = {};
const rejectHandler = config.rejectHandler || defaultRejectHandler;

self.jobQueue = new MiniQueue(config);
self.createMiniQueue = function (config) {
const miniQueue = new MiniQueue(config);

// When any task has gone to `process` state,
// wait for `end` event of `res` object (`job.data.res`)
// then leave the `process` state
miniQueue.on('process', function (job, done) {
//job.data.res.removeAllListeners('end'); // Remove listener which was set in on(queued) event
job.data.res.once('end', function () { // `end` event is sent on res.end() by `express-end` middleware package
done();
});
job.data.next();
});

// When any task has gone to `process` state,
// wait for `end` event of `res` object (`job.data.res`)
// then leave the `process` state
self.jobQueue.on('process', function(job, done) {
//job.data.res.removeAllListeners('end'); // Remove listener which was set in on(queued) event
job.data.res.once('end', function() { // `end` event is sent on res.end() by `express-end` middleware package
done();
miniQueue.on('reject', function (job) {
debug('Rejected ' + job.data.req.path);
rejectHandler(job.data.req, job.data.res);
});
job.data.next();
});

self.jobQueue.on('reject', function(job) {
debug('Rejected ' + job.data.req.path);
rejectHandler(job.data.req, job.data.res);
});

self.queueMw = function(req,res,next) {
// expose queue
resultMw.queue = miniQueue;

const data = { req: req, res: res, next: next };
const job = self.jobQueue.createJob(data);

// Set listeners for logging
res.once('close', function() { job.log('resOnClose'); }); // Closed from remote end
res.once('end', function() { job.log('resOnEnd'); }); // Sent on res.end() by express-end
res.once('finish', function() { job.log('resOnFinish'); });

// Handle disconnect from client while in queue
res.once('close', function() {
if (job.status === 'queue') {
// Return HTTP 204 No Content
// - it must not be sent as the connection is already close by Client
job.data.res.status(204).end();
self.jobQueue._cancelJob(job);
}
});
return miniQueue;

}

self.queueMw = function (req, res, next) {

const data = { req: req, res: res, next: next };
const fieldValue = config.field ? (req.body || {})[config.field] || (req.params || {})[config.field] || (req.query || {})[config.field] : null;

if (fieldValue) {
if (!self.jobQueue) self.jobQueue = {};
if (!self.jobQueue[fieldValue]) self.jobQueue[fieldValue] = self.createMiniQueue(config);

const job = self.jobQueue[fieldValue].createJob(data);

// Set listeners for logging
res.once('close', function () { job.log('resOnClose'); }); // Closed from remote end
res.once('end', function () { job.log('resOnEnd'); }); // Sent on res.end() by express-end
res.once('finish', function () { job.log('resOnFinish'); });

// Handle disconnect from client while in queue
res.once('close', function () {
if (job.status === 'queue') {
// Return HTTP 204 No Content
// - it must not be sent as the connection is already close by Client
job.data.res.status(204).end();
self.jobQueue[fieldValue]._cancelJob(job);
}
});

} else {

if (!self.jobQueue)
self.jobQueue = self.createMiniQueue(config);

const job = self.jobQueue.createJob(data);

// Set listeners for logging
res.once('close', function () { job.log('resOnClose'); }); // Closed from remote end
res.once('end', function () { job.log('resOnEnd'); }); // Sent on res.end() by express-end
res.once('finish', function () { job.log('resOnFinish'); });

// Handle disconnect from client while in queue
res.once('close', function () {
if (job.status === 'queue') {
// Return HTTP 204 No Content
// - it must not be sent as the connection is already close by Client
job.data.res.status(204).end();
self.jobQueue._cancelJob(job);
}
});
}

};

// merge `end` and `queue` middlewares
const resultMw = function(req, res, next) {
const resultMw = function (req, res, next) {
endMw(req, res, function () { // Inject res.end() handler to emit 'end' event
self.queueMw(req, res, next); // Use this middleware
});
};

// expose queue
resultMw.queue = self.jobQueue;

return resultMw;
};
Expand Down