queue.js 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294
  1. 'use strict';
  2. Object.defineProperty(exports, "__esModule", {
  3. value: true
  4. });
  5. exports.default = queue;
  6. var _onlyOnce = require('./onlyOnce.js');
  7. var _onlyOnce2 = _interopRequireDefault(_onlyOnce);
  8. var _setImmediate = require('./setImmediate.js');
  9. var _setImmediate2 = _interopRequireDefault(_setImmediate);
  10. var _DoublyLinkedList = require('./DoublyLinkedList.js');
  11. var _DoublyLinkedList2 = _interopRequireDefault(_DoublyLinkedList);
  12. var _wrapAsync = require('./wrapAsync.js');
  13. var _wrapAsync2 = _interopRequireDefault(_wrapAsync);
  14. function _interopRequireDefault(obj) { return obj && obj.__esModule ? obj : { default: obj }; }
  15. function queue(worker, concurrency, payload) {
  16. if (concurrency == null) {
  17. concurrency = 1;
  18. } else if (concurrency === 0) {
  19. throw new RangeError('Concurrency must not be zero');
  20. }
  21. var _worker = (0, _wrapAsync2.default)(worker);
  22. var numRunning = 0;
  23. var workersList = [];
  24. const events = {
  25. error: [],
  26. drain: [],
  27. saturated: [],
  28. unsaturated: [],
  29. empty: []
  30. };
  31. function on(event, handler) {
  32. events[event].push(handler);
  33. }
  34. function once(event, handler) {
  35. const handleAndRemove = (...args) => {
  36. off(event, handleAndRemove);
  37. handler(...args);
  38. };
  39. events[event].push(handleAndRemove);
  40. }
  41. function off(event, handler) {
  42. if (!event) return Object.keys(events).forEach(ev => events[ev] = []);
  43. if (!handler) return events[event] = [];
  44. events[event] = events[event].filter(ev => ev !== handler);
  45. }
  46. function trigger(event, ...args) {
  47. events[event].forEach(handler => handler(...args));
  48. }
  49. var processingScheduled = false;
  50. function _insert(data, insertAtFront, rejectOnError, callback) {
  51. if (callback != null && typeof callback !== 'function') {
  52. throw new Error('task callback must be a function');
  53. }
  54. q.started = true;
  55. var res, rej;
  56. function promiseCallback(err, ...args) {
  57. // we don't care about the error, let the global error handler
  58. // deal with it
  59. if (err) return rejectOnError ? rej(err) : res();
  60. if (args.length <= 1) return res(args[0]);
  61. res(args);
  62. }
  63. var item = q._createTaskItem(data, rejectOnError ? promiseCallback : callback || promiseCallback);
  64. if (insertAtFront) {
  65. q._tasks.unshift(item);
  66. } else {
  67. q._tasks.push(item);
  68. }
  69. if (!processingScheduled) {
  70. processingScheduled = true;
  71. (0, _setImmediate2.default)(() => {
  72. processingScheduled = false;
  73. q.process();
  74. });
  75. }
  76. if (rejectOnError || !callback) {
  77. return new Promise((resolve, reject) => {
  78. res = resolve;
  79. rej = reject;
  80. });
  81. }
  82. }
  83. function _createCB(tasks) {
  84. return function (err, ...args) {
  85. numRunning -= 1;
  86. for (var i = 0, l = tasks.length; i < l; i++) {
  87. var task = tasks[i];
  88. var index = workersList.indexOf(task);
  89. if (index === 0) {
  90. workersList.shift();
  91. } else if (index > 0) {
  92. workersList.splice(index, 1);
  93. }
  94. task.callback(err, ...args);
  95. if (err != null) {
  96. trigger('error', err, task.data);
  97. }
  98. }
  99. if (numRunning <= q.concurrency - q.buffer) {
  100. trigger('unsaturated');
  101. }
  102. if (q.idle()) {
  103. trigger('drain');
  104. }
  105. q.process();
  106. };
  107. }
  108. function _maybeDrain(data) {
  109. if (data.length === 0 && q.idle()) {
  110. // call drain immediately if there are no tasks
  111. (0, _setImmediate2.default)(() => trigger('drain'));
  112. return true;
  113. }
  114. return false;
  115. }
  116. const eventMethod = name => handler => {
  117. if (!handler) {
  118. return new Promise((resolve, reject) => {
  119. once(name, (err, data) => {
  120. if (err) return reject(err);
  121. resolve(data);
  122. });
  123. });
  124. }
  125. off(name);
  126. on(name, handler);
  127. };
  128. var isProcessing = false;
  129. var q = {
  130. _tasks: new _DoublyLinkedList2.default(),
  131. _createTaskItem(data, callback) {
  132. return {
  133. data,
  134. callback
  135. };
  136. },
  137. *[Symbol.iterator]() {
  138. yield* q._tasks[Symbol.iterator]();
  139. },
  140. concurrency,
  141. payload,
  142. buffer: concurrency / 4,
  143. started: false,
  144. paused: false,
  145. push(data, callback) {
  146. if (Array.isArray(data)) {
  147. if (_maybeDrain(data)) return;
  148. return data.map(datum => _insert(datum, false, false, callback));
  149. }
  150. return _insert(data, false, false, callback);
  151. },
  152. pushAsync(data, callback) {
  153. if (Array.isArray(data)) {
  154. if (_maybeDrain(data)) return;
  155. return data.map(datum => _insert(datum, false, true, callback));
  156. }
  157. return _insert(data, false, true, callback);
  158. },
  159. kill() {
  160. off();
  161. q._tasks.empty();
  162. },
  163. unshift(data, callback) {
  164. if (Array.isArray(data)) {
  165. if (_maybeDrain(data)) return;
  166. return data.map(datum => _insert(datum, true, false, callback));
  167. }
  168. return _insert(data, true, false, callback);
  169. },
  170. unshiftAsync(data, callback) {
  171. if (Array.isArray(data)) {
  172. if (_maybeDrain(data)) return;
  173. return data.map(datum => _insert(datum, true, true, callback));
  174. }
  175. return _insert(data, true, true, callback);
  176. },
  177. remove(testFn) {
  178. q._tasks.remove(testFn);
  179. },
  180. process() {
  181. // Avoid trying to start too many processing operations. This can occur
  182. // when callbacks resolve synchronously (#1267).
  183. if (isProcessing) {
  184. return;
  185. }
  186. isProcessing = true;
  187. while (!q.paused && numRunning < q.concurrency && q._tasks.length) {
  188. var tasks = [],
  189. data = [];
  190. var l = q._tasks.length;
  191. if (q.payload) l = Math.min(l, q.payload);
  192. for (var i = 0; i < l; i++) {
  193. var node = q._tasks.shift();
  194. tasks.push(node);
  195. workersList.push(node);
  196. data.push(node.data);
  197. }
  198. numRunning += 1;
  199. if (q._tasks.length === 0) {
  200. trigger('empty');
  201. }
  202. if (numRunning === q.concurrency) {
  203. trigger('saturated');
  204. }
  205. var cb = (0, _onlyOnce2.default)(_createCB(tasks));
  206. _worker(data, cb);
  207. }
  208. isProcessing = false;
  209. },
  210. length() {
  211. return q._tasks.length;
  212. },
  213. running() {
  214. return numRunning;
  215. },
  216. workersList() {
  217. return workersList;
  218. },
  219. idle() {
  220. return q._tasks.length + numRunning === 0;
  221. },
  222. pause() {
  223. q.paused = true;
  224. },
  225. resume() {
  226. if (q.paused === false) {
  227. return;
  228. }
  229. q.paused = false;
  230. (0, _setImmediate2.default)(q.process);
  231. }
  232. };
  233. // define these as fixed properties, so people get useful errors when updating
  234. Object.defineProperties(q, {
  235. saturated: {
  236. writable: false,
  237. value: eventMethod('saturated')
  238. },
  239. unsaturated: {
  240. writable: false,
  241. value: eventMethod('unsaturated')
  242. },
  243. empty: {
  244. writable: false,
  245. value: eventMethod('empty')
  246. },
  247. drain: {
  248. writable: false,
  249. value: eventMethod('drain')
  250. },
  251. error: {
  252. writable: false,
  253. value: eventMethod('error')
  254. }
  255. });
  256. return q;
  257. }
  258. module.exports = exports['default'];