source: frontend/node_modules/async/internal/queue.js

Last change on this file was 9af201e, checked in by MBK <marija.karapandzova@…>, 11 days ago

Fix frontend appearance

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