forked from singlesly/bingx-api
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathbingx-market-socket-stream.ts
More file actions
95 lines (84 loc) · 2.69 KB
/
Copy pathbingx-market-socket-stream.ts
File metadata and controls
95 lines (84 loc) · 2.69 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
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
import { pong } from 'bingx-api/bingx';
import { webSocket } from 'rxjs/webSocket';
import { BingxWebsocketDeserializer } from 'bingx-api/bingx-socket/bingx-websocket-deserializer';
import { BingxWebsocketSerializer } from 'bingx-api/bingx-socket/bingx-websocket-serializer';
import * as WebSocket from 'ws';
import { filterAndEmitToSubject } from 'bingx-api/bingx-socket/operators/filter-and-emit-to-subject';
import {
BehaviorSubject,
distinct,
ReplaySubject,
Subject,
switchMap,
tap,
} from 'rxjs';
import { HeartbeatInterface } from 'bingx-api/bingx-socket/interfaces/heartbeat.interface';
import {
LatestTradeEvent,
MarkerSubscription,
MarketWebsocketEvents,
SubscriptionType,
} from 'bingx-api/bingx-socket/events/market-websocket-events';
export class BingxMarketSocketStream {
private forceClose$ = new BehaviorSubject<boolean>(false);
private readonly dataTypes$ = new ReplaySubject<SubscriptionType>();
private readonly onConnect$ = new Subject();
public readonly onDisconnect$ = new Subject<CloseEvent>();
public readonly heartbeat$ = new ReplaySubject<HeartbeatInterface>(1);
public readonly latestTradeDetail$ = new Subject<LatestTradeEvent>();
constructor(
url: URL = new URL('/swap-market', 'wss://open-api-swap.bingx.com'),
) {
this.connect(url);
}
private async connect(url: URL): Promise<void> {
const socket$ = webSocket<MarketWebsocketEvents | MarkerSubscription>({
deserializer: (e) => new BingxWebsocketDeserializer().deserializer(e),
serializer: (e) => new BingxWebsocketSerializer().serializer(e),
url: url.toString(),
WebSocketCtor: WebSocket as never,
openObserver: this.onConnect$,
closeObserver: this.onDisconnect$,
});
const socketSubscription = socket$
.pipe(
pong(socket$, this.heartbeat$),
filterAndEmitToSubject(
(event): event is LatestTradeEvent =>
event.dataType.includes('trade'),
this.latestTradeDetail$,
),
)
.subscribe();
this.onConnect$
.pipe(
switchMap(() => this.dataTypes$),
distinct(),
tap((dataType) =>
socket$.next({
id: `listen-for-${dataType}`,
reqType: 'sub',
dataType,
}),
),
)
.subscribe();
this.onDisconnect$.subscribe(() => {
socketSubscription.unsubscribe();
if (!this.forceClose$.value) {
this.connect(url);
}
});
this.forceClose$.subscribe((v) => {
if (v) {
socket$.complete();
}
});
}
public disconnect() {
this.forceClose$.next(true);
}
public subscribe(dataType: SubscriptionType) {
this.dataTypes$.next(dataType);
}
}