/
SubscriptionIterator.ts
68 lines (60 loc) · 1.93 KB
/
SubscriptionIterator.ts
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
64
65
66
67
68
export type SubscriptionTransform = (value: any) => any
export type SubscriptionListener = (value: JSON) => Promise<void>
interface Handlers {
onStart?: (listener: SubscriptionListener) => void
onCompleted?: (listener: SubscriptionListener) => void
transform?: SubscriptionTransform
}
export default class SubscriptionIterator<T = any> implements AsyncIterator<T> {
private done = false
private pushQueue = [] as any[]
private pullQueue = [] as ((resolvedValue?: any) => void)[]
private transform?: SubscriptionTransform
private onCompleted?: (listener: SubscriptionListener) => void
private pushValue: SubscriptionListener = async (input) => {
const value = this.transform ? await this.transform(input) : input
if (value !== undefined) {
const resolver = this.pullQueue.shift()
if (resolver) {
resolver({value, done: false})
} else {
this.pushQueue.push(value)
}
}
}
constructor({onStart, onCompleted, transform}: Handlers) {
this.transform = transform
this.onCompleted = onCompleted
onStart?.(this.pushValue)
}
private close = () => {
if (this.done) return
this.done = true
this.onCompleted?.(this.pushValue)
this.pullQueue.forEach((resolve) => resolve({done: true, value: undefined}))
this.pullQueue = []
};
[Symbol.asyncIterator]() {
return this
}
next() {
return new Promise<IteratorResult<any>>((resolve) => {
if (this.done) resolve({done: true, value: undefined})
const value = this.pushQueue.shift()
if (value) {
resolve({value, done: false})
} else {
this.pullQueue.push(resolve)
}
})
}
return() {
this.close()
return Promise.resolve({done: true as const, value: undefined})
}
throw(error: unknown) {
const value = error instanceof Error ? error : new Error('SubscriptionIterator Error')
this.close()
return Promise.resolve({done: true as const, value})
}
}