groupBy.js
2.77 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
import { Observable } from '../Observable';
import { innerFrom } from '../observable/innerFrom';
import { Subject } from '../Subject';
import { operate } from '../util/lift';
import { createOperatorSubscriber, OperatorSubscriber } from './OperatorSubscriber';
export function groupBy(keySelector, elementOrOptions, duration, connector) {
return operate((source, subscriber) => {
let element;
if (!elementOrOptions || typeof elementOrOptions === 'function') {
element = elementOrOptions;
}
else {
({ duration, element, connector } = elementOrOptions);
}
const groups = new Map();
const notify = (cb) => {
groups.forEach(cb);
cb(subscriber);
};
const handleError = (err) => notify((consumer) => consumer.error(err));
let activeGroups = 0;
let teardownAttempted = false;
const groupBySourceSubscriber = new OperatorSubscriber(subscriber, (value) => {
try {
const key = keySelector(value);
let group = groups.get(key);
if (!group) {
groups.set(key, (group = connector ? connector() : new Subject()));
const grouped = createGroupedObservable(key, group);
subscriber.next(grouped);
if (duration) {
const durationSubscriber = createOperatorSubscriber(group, () => {
group.complete();
durationSubscriber === null || durationSubscriber === void 0 ? void 0 : durationSubscriber.unsubscribe();
}, undefined, undefined, () => groups.delete(key));
groupBySourceSubscriber.add(innerFrom(duration(grouped)).subscribe(durationSubscriber));
}
}
group.next(element ? element(value) : value);
}
catch (err) {
handleError(err);
}
}, () => notify((consumer) => consumer.complete()), handleError, () => groups.clear(), () => {
teardownAttempted = true;
return activeGroups === 0;
});
source.subscribe(groupBySourceSubscriber);
function createGroupedObservable(key, groupSubject) {
const result = new Observable((groupSubscriber) => {
activeGroups++;
const innerSub = groupSubject.subscribe(groupSubscriber);
return () => {
innerSub.unsubscribe();
--activeGroups === 0 && teardownAttempted && groupBySourceSubscriber.unsubscribe();
};
});
result.key = key;
return result;
}
});
}
//# sourceMappingURL=groupBy.js.map