-
Notifications
You must be signed in to change notification settings - Fork 11
/
through2-concurrent.js
101 lines (88 loc) · 2.78 KB
/
through2-concurrent.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
// Like through2 except execute in parallel with a set maximum
// concurrency
"use strict";
var through2 = require('through2');
function cbNoop (cb) {
cb();
}
module.exports = function concurrentThrough (options, transform, flush) {
var concurrent = 0, lastCallback = null, pendingFinish = null;
if (typeof options === 'function') {
flush = transform;
transform = options;
options = {};
}
var maxConcurrency = options.maxConcurrency || 16;
function _transform (message, enc, callback) {
var self = this;
var callbackCalled = false;
concurrent++;
if (concurrent < maxConcurrency) {
// Ask for more right away
callback();
} else {
// We're at the concurrency limit, save the callback for
// when we're ready for more
lastCallback = callback;
}
transform.call(this, message, enc, function (err) {
// Ignore multiple calls of the callback (shouldn't ever
// happen, but just in case)
if (callbackCalled) return;
callbackCalled = true;
if (err) {
self.emit('error', err);
} else if (arguments.length > 1) {
self.push(arguments[1]);
}
concurrent--;
if (lastCallback) {
var cb = lastCallback;
lastCallback = null;
cb();
}
if (concurrent === 0 && pendingFinish) {
pendingFinish();
pendingFinish = null;
}
});
}
// We need to pass in final to through2 even if the caller has
// not given us a final option so that it will wait for all
// transform callbacks to complete before emitting a "finish"
// and "end" event.
if (typeof options.final !== 'function') {
options.final = cbNoop;
}
// We also wrap flush to make sure anyone using an ancient version
// of through2 without support for final will get the old behaviour.
// TODO: don't wrap flush after upgrading through2 to a version with guaranteed `_final`
if (typeof flush !== 'function') {
flush = cbNoop;
}
// Flush is always called only after Final has finished
// to ensure that data from Final gets processed, so we only need one pending callback at a time
function callOnFinish (original) {
return function (callback) {
if (concurrent === 0) {
original.call(this, callback);
} else {
pendingFinish = original.bind(this, callback);
}
}
}
options.final = callOnFinish(options.final);
return through2(options, _transform, callOnFinish(flush));
};
module.exports.obj = function (options, transform, flush) {
if (typeof options === 'function') {
flush = transform;
transform = options;
options = {};
}
options.objectMode = true;
if (options.highWaterMark == null) {
options.highWaterMark = 16;
}
return module.exports(options, transform, flush);
};