1 | import { identity } from '../util/identity';
|
2 | import { wrapWithAbort } from './operators/withabort';
|
3 | import { safeRace } from '../util/safeRace';
|
4 |
|
5 | const NEVER_PROMISE = new Promise(() => { });
|
6 | function wrapPromiseWithIndex(promise, index) {
|
7 | return promise.then((value) => ({ value, index }));
|
8 | }
|
9 |
|
10 |
|
11 |
|
12 |
|
13 |
|
14 |
|
15 |
|
16 |
|
17 | export async function forkJoin(...sources) {
|
18 | let signal = sources.shift();
|
19 | if (!(signal instanceof AbortSignal)) {
|
20 | sources.unshift(signal);
|
21 | signal = undefined;
|
22 | }
|
23 | const length = sources.length;
|
24 | const iterators = new Array(length);
|
25 | const nexts = new Array(length);
|
26 | let active = length;
|
27 | const values = new Array(length);
|
28 | const hasValues = new Array(length);
|
29 | hasValues.fill(false);
|
30 | for (let i = 0; i < length; i++) {
|
31 | const iterator = wrapWithAbort(sources[i], signal)[Symbol.asyncIterator]();
|
32 | iterators[i] = iterator;
|
33 | nexts[i] = wrapPromiseWithIndex(iterator.next(), i);
|
34 | }
|
35 | while (active > 0) {
|
36 | const next = safeRace(nexts);
|
37 | const { value: next$, index } = await next;
|
38 | if (next$.done) {
|
39 | nexts[index] = NEVER_PROMISE;
|
40 | active--;
|
41 | }
|
42 | else {
|
43 | const iterator$ = iterators[index];
|
44 | nexts[index] = wrapPromiseWithIndex(iterator$.next(), index);
|
45 | hasValues[index] = true;
|
46 | values[index] = next$.value;
|
47 | }
|
48 | }
|
49 | if (hasValues.length > 0 && hasValues.every(identity)) {
|
50 | return values;
|
51 | }
|
52 | return undefined;
|
53 | }
|
54 |
|
55 |
|