groupBy.js 3.2 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667
  1. "use strict";
  2. Object.defineProperty(exports, "__esModule", { value: true });
  3. exports.groupBy = void 0;
  4. var Observable_1 = require("../Observable");
  5. var innerFrom_1 = require("../observable/innerFrom");
  6. var Subject_1 = require("../Subject");
  7. var lift_1 = require("../util/lift");
  8. var OperatorSubscriber_1 = require("./OperatorSubscriber");
  9. function groupBy(keySelector, elementOrOptions, duration, connector) {
  10. return lift_1.operate(function (source, subscriber) {
  11. var element;
  12. if (!elementOrOptions || typeof elementOrOptions === 'function') {
  13. element = elementOrOptions;
  14. }
  15. else {
  16. (duration = elementOrOptions.duration, element = elementOrOptions.element, connector = elementOrOptions.connector);
  17. }
  18. var groups = new Map();
  19. var notify = function (cb) {
  20. groups.forEach(cb);
  21. cb(subscriber);
  22. };
  23. var handleError = function (err) { return notify(function (consumer) { return consumer.error(err); }); };
  24. var activeGroups = 0;
  25. var teardownAttempted = false;
  26. var groupBySourceSubscriber = new OperatorSubscriber_1.OperatorSubscriber(subscriber, function (value) {
  27. try {
  28. var key_1 = keySelector(value);
  29. var group_1 = groups.get(key_1);
  30. if (!group_1) {
  31. groups.set(key_1, (group_1 = connector ? connector() : new Subject_1.Subject()));
  32. var grouped = createGroupedObservable(key_1, group_1);
  33. subscriber.next(grouped);
  34. if (duration) {
  35. var durationSubscriber_1 = OperatorSubscriber_1.createOperatorSubscriber(group_1, function () {
  36. group_1.complete();
  37. durationSubscriber_1 === null || durationSubscriber_1 === void 0 ? void 0 : durationSubscriber_1.unsubscribe();
  38. }, undefined, undefined, function () { return groups.delete(key_1); });
  39. groupBySourceSubscriber.add(innerFrom_1.innerFrom(duration(grouped)).subscribe(durationSubscriber_1));
  40. }
  41. }
  42. group_1.next(element ? element(value) : value);
  43. }
  44. catch (err) {
  45. handleError(err);
  46. }
  47. }, function () { return notify(function (consumer) { return consumer.complete(); }); }, handleError, function () { return groups.clear(); }, function () {
  48. teardownAttempted = true;
  49. return activeGroups === 0;
  50. });
  51. source.subscribe(groupBySourceSubscriber);
  52. function createGroupedObservable(key, groupSubject) {
  53. var result = new Observable_1.Observable(function (groupSubscriber) {
  54. activeGroups++;
  55. var innerSub = groupSubject.subscribe(groupSubscriber);
  56. return function () {
  57. innerSub.unsubscribe();
  58. --activeGroups === 0 && teardownAttempted && groupBySourceSubscriber.unsubscribe();
  59. };
  60. });
  61. result.key = key;
  62. return result;
  63. }
  64. });
  65. }
  66. exports.groupBy = groupBy;
  67. //# sourceMappingURL=groupBy.js.map