This repository has been archived by the owner on Nov 29, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 2
/
GQLSubscription.ts
130 lines (120 loc) · 3.14 KB
/
GQLSubscription.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
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
import { FeatureRunner } from '../runner'
import { expect } from 'chai'
import * as Websocket from 'ws'
import { Client, Message } from 'paho-mqtt'
// @ts-ignore
global.WebSocket = Websocket // required for Paho MQTT
export class GQLSubscription {
public readonly connection: Promise<Client>
public messages: any[] = []
private readonly client: Client
private readonly subscribers: {
id: string
matches: (msg: { [key: string]: string }) => boolean
onMatch: ((msg: any) => void)[]
}[] = []
private readonly subscriberMessages: { [key: string]: any[] } = {}
constructor(
selection: string,
url: string,
clientId: string,
topics: string[],
runner: FeatureRunner<any>,
) {
this.client = new Client(url, clientId)
this.client.onMessageArrived = (msg: Message) => {
const {
data: { [selection]: result },
} = JSON.parse(msg.payloadString)
this.messages.push(result)
void runner.progress('<GQL@', msg.payloadString)
this.notifySubcribers(result)
}
this.connection = new Promise((resolve, reject) => {
this.client.connect({
useSSL: false,
mqttVersion: 3,
onSuccess: async () => {
await Promise.all(
topics.map(
(topic) =>
new Promise((resolve1, reject1) => {
this.client.subscribe(topic, {
onSuccess: resolve1,
onFailure: reject1,
})
}),
),
)
resolve(this.client)
},
onFailure: reject,
})
})
}
notifySubcribers = (result: { [key: string]: string }): void => {
this.subscribers.forEach(({ id, matches, onMatch }) => {
if (matches(result)) {
if (this.subscriberMessages[id] !== undefined) {
this.subscriberMessages[id] = []
}
this.subscriberMessages[id].push(result)
onMatch.forEach((fn) => fn(result))
}
})
}
addListener = (
listenerId: string,
matcher: { [key: string]: string },
): void => {
this.subscribers.push({
id: listenerId,
matches: (message: any): boolean => {
try {
expect(message).to.containSubset(matcher)
return true
} catch (error) {
return false
}
},
onMatch: [],
})
// Notify about existing messages
this.messages.forEach((message) => this.notifySubcribers(message))
}
disconnect = (): void => {
this.client.disconnect()
}
/**
* Returns a message for the given subscription id within a certain time
*/
listenerMessage = async (
listenerId: string,
timeout = 5000,
): Promise<any> => {
if (
Array.isArray(this.subscriberMessages[listenerId]) &&
this.subscriberMessages[listenerId].length > 0
) {
return this.subscriberMessages[listenerId][
this.subscriberMessages[listenerId].length - 1
]
}
let messageListenerId: number
return new Promise<any>((resolve, reject) => {
const sub = this.subscribers.find(({ id }) => id === listenerId)
if (!sub) {
throw new Error(`Subscriber for "${listenerId}" not found!`)
}
const timeoutId = setTimeout(() => {
sub.onMatch.splice(messageListenerId, 1)
reject()
}, timeout)
// Register listener for arriving messages
messageListenerId = sub.onMatch.push((msg: any) => {
clearTimeout(timeoutId)
resolve(msg)
})
})
}
}