-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathindex.js
More file actions
109 lines (104 loc) · 3.27 KB
/
Copy pathindex.js
File metadata and controls
109 lines (104 loc) · 3.27 KB
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
102
103
104
105
106
107
108
109
"use strict";
var co = require('co');
module.exports.mapLimit = function (array, numberOfProcesses, generator) {
/** If wrong number of processes reject it */
if (numberOfProcesses < 1) {
return Promise.reject(new Error('Number of processes has to be positive number'));
}
/** If not array or empty array resolve it */
if (!Array.isArray(array)) {
return Promise.reject(new Error('No array passed'));
}
if (array.length === 0) {
return Promise.resolve(array);
}
var results = [], runningProcesses = 0, keys = array.map((item, key) => {
return key;
});
return new Promise(
/** recursive function inside Promise */
function runProcess(resolve, reject) {
if (keys.length > 0) {
/** asynchronous iteration over maximum empty workers */
keys.splice(0, numberOfProcesses - runningProcesses).forEach((key) => {
runningProcesses++;
/** run generator */
co(generator(array[key], key, array))
.then((res) => {
/** save generator's result to results array under it's original index */
results[key] = res;
/** decrement number of running processes */
runningProcesses--;
runProcess(resolve, reject)
})
.catch((rej) => {
results[key] = {error: rej};
runningProcesses--;
runProcess(resolve, reject)
})
})
} else if (runningProcesses === 0 && keys.length === 0) {
/** if there is no more running process and it has already iterated array length
* then resolve it */
resolve(results);
}
}
);
};
/**
*
* @param array
* @param delay Delay in miliseconds
* @param generator
* @returns {Promise}
*/
module.exports.mapDelay = function (array, delay, generator) {
/** If wrong number of processes reject it */
if (delay < 0) {
return Promise.reject(new Error('Delay has to be positive number or zero'));
}
/** If not array or empty array resolve it */
if (!Array.isArray(array)) {
return Promise.reject(new Error('No array passed'));
}
if (array.length === 0) {
return Promise.resolve(array);
}
var results = [], keys = array.map((item, key) => {
return key;
});
return new Promise(
/** recursive function inside Promise */
function runProcess(resolve, reject) {
if (keys.length > 0) {
/** run generator */
let key = keys.splice(0, 1)[0];
co(generator(array[key], key, array))
.then((res) => {
/** save generator's result to results array under it's original index */
results[key] = res;
/** if it has already iterated array length
* then resolve it
* or recurse */
if (keys.length === 0) {
resolve(results)
} else {
setTimeout(function () {
runProcess(resolve, reject)
}, delay)
}
})
.catch((rej) => {
results[key] = {error: rej};
if (keys.length === 0) {
resolve(results)
} else {
setTimeout(function () {
runProcess(resolve, reject)
}, delay)
}
})
}
}
);
};