1 | "use strict";
|
2 | var __awaiter = (this && this.__awaiter) || function (thisArg, _arguments, P, generator) {
|
3 | function adopt(value) { return value instanceof P ? value : new P(function (resolve) { resolve(value); }); }
|
4 | return new (P || (P = Promise))(function (resolve, reject) {
|
5 | function fulfilled(value) { try { step(generator.next(value)); } catch (e) { reject(e); } }
|
6 | function rejected(value) { try { step(generator["throw"](value)); } catch (e) { reject(e); } }
|
7 | function step(result) { result.done ? resolve(result.value) : adopt(result.value).then(fulfilled, rejected); }
|
8 | step((generator = generator.apply(thisArg, _arguments || [])).next());
|
9 | });
|
10 | };
|
11 | Object.defineProperty(exports, "__esModule", { value: true });
|
12 | const iterall_1 = require("iterall");
|
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 | class PubSubAsyncIterator {
|
43 | constructor(pubsub, eventNames) {
|
44 | this.pubsub = pubsub;
|
45 | this.pullQueue = [];
|
46 | this.pushQueue = [];
|
47 | this.listening = true;
|
48 | this.eventsArray =
|
49 | typeof eventNames === 'string' ? [eventNames] : eventNames;
|
50 | this.allSubscribed = this.subscribeAll();
|
51 | }
|
52 | next() {
|
53 | return __awaiter(this, void 0, void 0, function* () {
|
54 | yield this.allSubscribed;
|
55 | return this.listening ? this.pullValue() : this.return();
|
56 | });
|
57 | }
|
58 | return() {
|
59 | return __awaiter(this, void 0, void 0, function* () {
|
60 | this.emptyQueue(yield this.allSubscribed);
|
61 | return { value: undefined, done: true };
|
62 | });
|
63 | }
|
64 | throw(error) {
|
65 | return __awaiter(this, void 0, void 0, function* () {
|
66 | this.emptyQueue(yield this.allSubscribed);
|
67 | return Promise.reject(error);
|
68 | });
|
69 | }
|
70 | [iterall_1.$$asyncIterator]() {
|
71 | return this;
|
72 | }
|
73 | pushValue(event) {
|
74 | return __awaiter(this, void 0, void 0, function* () {
|
75 | yield this.allSubscribed;
|
76 | if (this.pullQueue.length !== 0) {
|
77 | this.pullQueue.shift()({ value: event, done: false });
|
78 | }
|
79 | else {
|
80 | this.pushQueue.push(event);
|
81 | }
|
82 | });
|
83 | }
|
84 | pullValue() {
|
85 | return new Promise(resolve => {
|
86 | if (this.pushQueue.length !== 0) {
|
87 | resolve({ value: this.pushQueue.shift(), done: false });
|
88 | }
|
89 | else {
|
90 | this.pullQueue.push(resolve);
|
91 | }
|
92 | });
|
93 | }
|
94 | emptyQueue(subscriptionIds) {
|
95 | if (this.listening) {
|
96 | this.listening = false;
|
97 | this.unsubscribeAll(subscriptionIds);
|
98 | this.pullQueue.forEach(resolve => resolve({ value: undefined, done: true }));
|
99 | this.pullQueue.length = 0;
|
100 | this.pushQueue.length = 0;
|
101 | }
|
102 | }
|
103 | subscribeAll() {
|
104 | return Promise.all(this.eventsArray.map(eventName => this.pubsub.subscribe(eventName, this.pushValue.bind(this), {})));
|
105 | }
|
106 | unsubscribeAll(subscriptionIds) {
|
107 | for (const subscriptionId of subscriptionIds) {
|
108 | this.pubsub.unsubscribe(subscriptionId);
|
109 | }
|
110 | }
|
111 | }
|
112 | exports.PubSubAsyncIterator = PubSubAsyncIterator;
|