forked from singlesly/bingx-api
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathbingx-account-socket-stream.ts
More file actions
122 lines (110 loc) · 4.08 KB
/
Copy pathbingx-account-socket-stream.ts
File metadata and controls
122 lines (110 loc) · 4.08 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
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
import { BehaviorSubject, ReplaySubject, Subject } from 'rxjs';
import { AccountInterface } from 'bingx-api/bingx/account/account.interface';
import {
BingxGenerateListenKeyEndpoint,
HttpRequestExecutor,
pong,
RequestExecutorInterface,
} 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 { HeartbeatInterface } from 'bingx-api/bingx-socket/interfaces/heartbeat.interface';
import { filterAndEmitToSubject } from 'bingx-api/bingx-socket/operators/filter-and-emit-to-subject';
import {
AccountBalanceAndPositionPushEvent,
AccountOrderUpdatePushEvent,
AccountWebSocketEvent,
AccountWebsocketEventType,
ListenKeyExpiredEvent,
} from 'bingx-api/bingx-socket/events/account-websocket-events';
import * as WebSocket from 'ws';
export interface BingxAccountSocketStreamConfiguration {
requestExecutor?: RequestExecutorInterface;
url?: URL;
}
export class BingxAccountSocketStream {
private forceClose$ = new BehaviorSubject<boolean>(false);
private readonly configuration: Required<BingxAccountSocketStreamConfiguration>;
public readonly onConnect$ = new Subject();
public readonly onDisconnect$ = new Subject<CloseEvent>();
public readonly heartbeat$ = new ReplaySubject<HeartbeatInterface>(1);
public readonly listenKeyExpired$ = new Subject<ListenKeyExpiredEvent>();
public readonly accountBalanceAndPositionPush$ =
new Subject<AccountBalanceAndPositionPushEvent>();
public readonly accountOrderUpdatePushEvent$ =
new Subject<AccountOrderUpdatePushEvent>();
constructor(
private readonly account: AccountInterface,
configuration: BingxAccountSocketStreamConfiguration = {},
) {
this.configuration = {
requestExecutor:
configuration.requestExecutor ?? new HttpRequestExecutor(),
url:
configuration.url ??
new URL('/swap-market', 'wss://open-api-swap.bingx.com'),
};
try {
this.connect(this.account, this.configuration.requestExecutor).catch(
() => {
console.error('Cannot connect to account', this.account);
},
);
} catch (e) {
console.error('Cannot connect to account error', this.account);
}
}
private async connect(
account: AccountInterface,
requestExecutor: RequestExecutorInterface,
): Promise<void> {
const responseKey = await requestExecutor.execute(
new BingxGenerateListenKeyEndpoint(account),
);
const url = this.configuration.url;
url.searchParams.set('listenKey', responseKey.listenKey);
const socket$ = webSocket<AccountWebSocketEvent>({
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 ListenKeyExpiredEvent =>
event.e === AccountWebsocketEventType.LISTEN_KEY_EXPIRED,
this.listenKeyExpired$,
),
filterAndEmitToSubject(
(event): event is AccountBalanceAndPositionPushEvent =>
event.e === AccountWebsocketEventType.ACCOUNT_UPDATE,
this.accountBalanceAndPositionPush$,
),
filterAndEmitToSubject(
(event): event is AccountOrderUpdatePushEvent =>
event.e === AccountWebsocketEventType.ORDER_TRADE_UPDATE,
this.accountOrderUpdatePushEvent$,
),
)
.subscribe();
this.onDisconnect$.subscribe(() => {
socketSubscription.unsubscribe();
if (!this.forceClose$.value) {
this.connect(account, requestExecutor);
}
});
this.forceClose$.subscribe((v) => {
if (v) {
socket$.complete();
}
});
}
public disconnect() {
this.forceClose$.next(true);
}
}