| 1 | 'use strict';
|
|---|
| 2 |
|
|---|
| 3 | Object.defineProperty(exports, "__esModule", {
|
|---|
| 4 | value: true
|
|---|
| 5 | });
|
|---|
| 6 |
|
|---|
| 7 | var _once = require('./once.js');
|
|---|
| 8 |
|
|---|
| 9 | var _once2 = _interopRequireDefault(_once);
|
|---|
| 10 |
|
|---|
| 11 | var _iterator = require('./iterator.js');
|
|---|
| 12 |
|
|---|
| 13 | var _iterator2 = _interopRequireDefault(_iterator);
|
|---|
| 14 |
|
|---|
| 15 | var _onlyOnce = require('./onlyOnce.js');
|
|---|
| 16 |
|
|---|
| 17 | var _onlyOnce2 = _interopRequireDefault(_onlyOnce);
|
|---|
| 18 |
|
|---|
| 19 | var _wrapAsync = require('./wrapAsync.js');
|
|---|
| 20 |
|
|---|
| 21 | var _asyncEachOfLimit = require('./asyncEachOfLimit.js');
|
|---|
| 22 |
|
|---|
| 23 | var _asyncEachOfLimit2 = _interopRequireDefault(_asyncEachOfLimit);
|
|---|
| 24 |
|
|---|
| 25 | var _breakLoop = require('./breakLoop.js');
|
|---|
| 26 |
|
|---|
| 27 | var _breakLoop2 = _interopRequireDefault(_breakLoop);
|
|---|
| 28 |
|
|---|
| 29 | function _interopRequireDefault(obj) { return obj && obj.__esModule ? obj : { default: obj }; }
|
|---|
| 30 |
|
|---|
| 31 | exports.default = limit => {
|
|---|
| 32 | return (obj, iteratee, callback) => {
|
|---|
| 33 | callback = (0, _once2.default)(callback);
|
|---|
| 34 | if (limit <= 0) {
|
|---|
| 35 | throw new RangeError('concurrency limit cannot be less than 1');
|
|---|
| 36 | }
|
|---|
| 37 | if (!obj) {
|
|---|
| 38 | return callback(null);
|
|---|
| 39 | }
|
|---|
| 40 | if ((0, _wrapAsync.isAsyncGenerator)(obj)) {
|
|---|
| 41 | return (0, _asyncEachOfLimit2.default)(obj, limit, iteratee, callback);
|
|---|
| 42 | }
|
|---|
| 43 | if ((0, _wrapAsync.isAsyncIterable)(obj)) {
|
|---|
| 44 | return (0, _asyncEachOfLimit2.default)(obj[Symbol.asyncIterator](), limit, iteratee, callback);
|
|---|
| 45 | }
|
|---|
| 46 | var nextElem = (0, _iterator2.default)(obj);
|
|---|
| 47 | var done = false;
|
|---|
| 48 | var canceled = false;
|
|---|
| 49 | var running = 0;
|
|---|
| 50 | var looping = false;
|
|---|
| 51 |
|
|---|
| 52 | function iterateeCallback(err, value) {
|
|---|
| 53 | if (canceled) return;
|
|---|
| 54 | running -= 1;
|
|---|
| 55 | if (err) {
|
|---|
| 56 | done = true;
|
|---|
| 57 | callback(err);
|
|---|
| 58 | } else if (err === false) {
|
|---|
| 59 | done = true;
|
|---|
| 60 | canceled = true;
|
|---|
| 61 | } else if (value === _breakLoop2.default || done && running <= 0) {
|
|---|
| 62 | done = true;
|
|---|
| 63 | return callback(null);
|
|---|
| 64 | } else if (!looping) {
|
|---|
| 65 | replenish();
|
|---|
| 66 | }
|
|---|
| 67 | }
|
|---|
| 68 |
|
|---|
| 69 | function replenish() {
|
|---|
| 70 | looping = true;
|
|---|
| 71 | while (running < limit && !done) {
|
|---|
| 72 | var elem = nextElem();
|
|---|
| 73 | if (elem === null) {
|
|---|
| 74 | done = true;
|
|---|
| 75 | if (running <= 0) {
|
|---|
| 76 | callback(null);
|
|---|
| 77 | }
|
|---|
| 78 | return;
|
|---|
| 79 | }
|
|---|
| 80 | running += 1;
|
|---|
| 81 | iteratee(elem.value, elem.key, (0, _onlyOnce2.default)(iterateeCallback));
|
|---|
| 82 | }
|
|---|
| 83 | looping = false;
|
|---|
| 84 | }
|
|---|
| 85 |
|
|---|
| 86 | replenish();
|
|---|
| 87 | };
|
|---|
| 88 | };
|
|---|
| 89 |
|
|---|
| 90 | module.exports = exports.default; |
|---|