Skip to content

Commit 2cd2c5b

Browse files
committed
add support for custom transport
1 parent 3ea2d0b commit 2cd2c5b

2 files changed

Lines changed: 165 additions & 103 deletions

File tree

lib/joiner.ts

Lines changed: 32 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import {Joiner, ClientAuthenticator, Serializer, JSONSerializer} from 'wampproto';
22

3-
import {BaseSession} from './types';
3+
import {BaseSession, Peer, WebSocketPeer} from './types';
44
import {getSubProtocol} from './helpers';
55

66

@@ -21,47 +21,45 @@ export class WAMPSessionJoiner {
2121
}
2222

2323
async join(uri: string, realm: string): Promise<BaseSession> {
24-
await ensureGlobalWebSocket()
24+
await ensureGlobalWebSocket();
2525
const ws = new WebSocket(uri, [getSubProtocol(this._serializer)]);
26+
await new Promise((resolve, reject) => {
27+
ws.addEventListener("open", resolve);
28+
ws.addEventListener("error", reject);
29+
});
30+
31+
const peer = new WebSocketPeer(ws);
32+
return joinPeer(peer, realm, this._serializer, this._authenticator);
33+
}
34+
}
2635

27-
const joiner = new Joiner(realm, this._serializer, this._authenticator);
2836

29-
ws.addEventListener('open', () => {
30-
ws.send(joiner.sendHello());
31-
});
37+
export async function joinPeer(
38+
peer: Peer,
39+
realm: string,
40+
serializer: Serializer,
41+
authenticator?: ClientAuthenticator
42+
): Promise<BaseSession> {
3243

33-
return new Promise<BaseSession>((resolve, reject) => {
34-
const wsMessageHandler = async (event: MessageEvent) => {
35-
try {
36-
let data = event.data;
44+
const joiner = new Joiner(realm, serializer, authenticator);
3745

38-
if (event.data instanceof Blob) {
39-
data = new Uint8Array(await event.data.arrayBuffer());
40-
}
46+
const hello = joiner.sendHello();
47+
peer.send(hello as Uint8Array);
4148

42-
const toSend = await joiner.receive(data);
43-
if (!toSend) {
44-
ws.removeEventListener('message', wsMessageHandler);
45-
ws.removeEventListener('close', closeHandler);
49+
// eslint-disable-next-line no-constant-condition
50+
while (true) {
51+
const msgBytes = await peer.receive();
4652

47-
const baseSession = new BaseSession(ws, wsMessageHandler, joiner.getSessionDetails(), this._serializer);
48-
resolve(baseSession);
49-
} else {
50-
ws.send(toSend);
51-
}
52-
} catch (error) {
53-
reject(error);
54-
}
55-
};
53+
const toSend = await joiner.receive(msgBytes);
5654

57-
const closeHandler = () => {
58-
ws.removeEventListener('message', wsMessageHandler);
59-
reject(new Error('Connection closed before handshake completed'));
60-
};
55+
if (toSend === null) {
56+
return new BaseSession(
57+
peer,
58+
joiner.getSessionDetails(),
59+
serializer
60+
);
61+
}
6162

62-
ws.addEventListener('message', wsMessageHandler);
63-
ws.addEventListener('error', (error) => reject(error));
64-
ws.addEventListener('close', closeHandler);
65-
});
63+
peer.send(toSend as Uint8Array);
6664
}
6765
}

lib/types.ts

Lines changed: 133 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -59,125 +59,87 @@ export abstract class IBaseSession {
5959
}
6060

6161
export class BaseSession extends IBaseSession {
62-
private readonly _ws: WebSocket;
63-
private readonly _wsMessageHandler: any;
64-
private readonly sessionDetails: SessionDetails;
62+
private readonly _peer: Peer;
6563
private readonly _serializer: Serializer;
64+
private readonly _sessionDetails: SessionDetails;
65+
6666
private _disconnectCallbacks: Array<(reason?: string) => Promise<void>> = [];
67-
private _queue: any[] = [];
68-
private _waiting: {resolve: (value: any) => void; reject: (reason?: any) => void;}[] = [];
69-
private _wsClosed = false;
7067

7168
constructor(
72-
ws: WebSocket,
73-
wsMessageHandler: any,
69+
peer: Peer,
7470
sessionDetails: SessionDetails,
7571
serializer: Serializer
7672
) {
7773
super();
78-
this._ws = ws;
79-
this._wsMessageHandler = wsMessageHandler;
80-
this.sessionDetails = sessionDetails;
74+
this._peer = peer;
8175
this._serializer = serializer;
82-
83-
this._ws.binaryType = "arraybuffer";
84-
this._ws.addEventListener("message", (event: MessageEvent) => {
85-
const data = event.data instanceof ArrayBuffer
86-
? new Uint8Array(event.data)
87-
: event.data;
88-
89-
if (this._waiting.length > 0) {
90-
const waiter = this._waiting.shift()!;
91-
waiter.resolve(data);
92-
} else {
93-
this._queue.push(data);
94-
}
95-
});
96-
97-
// close cleanly on abrupt client disconnect
98-
this._ws.addEventListener("close", async () => {
99-
this._wsClosed = true;
100-
101-
while (this._waiting.length > 0) {
102-
const waiter = this._waiting.shift()!;
103-
waiter.reject(new SessionClosedError());
104-
}
105-
106-
if (this._disconnectCallbacks.length > 0) {
107-
await Promise.all(this._disconnectCallbacks.map(cb => cb()));
108-
}
109-
await this.close();
110-
});
76+
this._sessionDetails = sessionDetails;
77+
78+
if (peer.onDisconnect) {
79+
peer.onDisconnect(async () => {
80+
if (this._disconnectCallbacks.length > 0) {
81+
await Promise.all(
82+
this._disconnectCallbacks.map(cb => cb())
83+
);
84+
}
85+
});
86+
}
11187
}
11288

11389
id(): number {
114-
return this.sessionDetails.sessionID;
90+
return this._sessionDetails.sessionID;
11591
}
11692

11793
realm(): string {
118-
return this.sessionDetails.realm;
94+
return this._sessionDetails.realm;
11995
}
12096

12197
authid(): string {
122-
return this.sessionDetails.authid;
98+
return this._sessionDetails.authid;
12399
}
124100

125101
authrole(): string {
126-
return this.sessionDetails.authrole;
102+
return this._sessionDetails.authrole;
127103
}
128104

129105
serializer(): Serializer {
130106
return this._serializer;
131107
}
132108

133109
send(data: any): void {
134-
this._ws.send(data);
135-
}
136-
137-
sendMessage(msg: Message): void {
138-
this.send(this._serializer.serialize(msg));
110+
this._peer.send(data);
139111
}
140112

141113
async receive(): Promise<any> {
142-
if (this._wsClosed) {
143-
throw new Error("Session closed");
144-
}
145-
146-
if (this._queue.length > 0) {
147-
return this._queue.shift();
148-
}
114+
return this._peer.receive();
115+
}
149116

150-
return new Promise((resolve, reject) => {
151-
this._waiting.push({ resolve, reject });
152-
});
117+
sendMessage(msg: Message): void {
118+
const bytes = this._serializer.serialize(msg);
119+
this._peer.send(bytes as Uint8Array);
153120
}
154121

155122
async receiveMessage(): Promise<Message> {
156-
return this._serializer.deserialize(await this.receive());
123+
const data = await this._peer.receive();
124+
return this._serializer.deserialize(data);
157125
}
158126

159127
async close(): Promise<void> {
160-
if (this._wsMessageHandler) {
161-
this._ws.removeEventListener("message", this._wsMessageHandler);
162-
this._ws.removeEventListener("close", this._wsMessageHandler);
163-
}
164-
165-
this._ws.close();
128+
await this._peer.close();
166129
}
167130

168131
isConnected(): boolean {
169-
return this._ws.readyState === WebSocket.OPEN;
132+
return this._peer.isConnected();
170133
}
171134

172135
onDisconnect(callback: (reason?: string) => Promise<void>): void {
173136
this._disconnectCallbacks.push(callback);
174137
}
175138

176139
getSessionDetails(): SessionDetails {
177-
return this.sessionDetails;
140+
return this._sessionDetails;
178141
}
179142
}
180-
181143
export class Result {
182144
args: any[];
183145
kwargs: { [key: string]: any };
@@ -315,3 +277,105 @@ export class Progress {
315277
}
316278

317279
export class SessionClosedError extends Error {}
280+
281+
export abstract class Peer {
282+
abstract send(data: Uint8Array): void;
283+
abstract receive(): Promise<Uint8Array>;
284+
abstract close(): Promise<void>;
285+
abstract isConnected(): boolean;
286+
onDisconnect?(callback: () => Promise<void>): void;
287+
}
288+
289+
export class WebSocketPeer extends Peer {
290+
private readonly _ws: WebSocket;
291+
private _queue: Uint8Array[] = [];
292+
private _waiting: {
293+
resolve: (data: Uint8Array) => void,
294+
reject: (err: Error) => void
295+
}[] = [];
296+
297+
private _disconnectHandlers: (() => Promise<void>)[] = [];
298+
private _closed = false;
299+
300+
constructor(ws: WebSocket) {
301+
super();
302+
303+
this._ws = ws;
304+
this._ws.binaryType = "arraybuffer";
305+
this._bindEvents();
306+
}
307+
308+
private _bindEvents() {
309+
this._ws.addEventListener("message", (event: MessageEvent) => {
310+
const data =
311+
event.data instanceof ArrayBuffer
312+
? new Uint8Array(event.data)
313+
: new Uint8Array(event.data);
314+
315+
if (this._waiting.length > 0) {
316+
const waiter = this._waiting.shift()!;
317+
waiter.resolve(data);
318+
} else {
319+
this._queue.push(data);
320+
}
321+
322+
});
323+
324+
this._ws.addEventListener("close", () => {
325+
this._handleDisconnect();
326+
});
327+
328+
this._ws.addEventListener("error", () => {
329+
this._handleDisconnect();
330+
});
331+
}
332+
333+
private async _handleDisconnect() {
334+
if (this._closed) return;
335+
this._closed = true;
336+
337+
while (this._waiting.length > 0) {
338+
const waiter = this._waiting.shift();
339+
waiter?.reject(new Error("WebSocket closed"));
340+
}
341+
342+
for (const cb of this._disconnectHandlers) {
343+
await cb();
344+
}
345+
}
346+
347+
onDisconnect(callback: () => Promise<void>): void {
348+
this._disconnectHandlers.push(callback);
349+
}
350+
351+
isConnected(): boolean {
352+
return this._ws.readyState === WebSocket.OPEN;
353+
}
354+
355+
send(data: Uint8Array): void {
356+
this._ws.send(data);
357+
}
358+
359+
async receive(): Promise<Uint8Array> {
360+
if (this._closed) {
361+
throw new Error("WebSocket closed in receive");
362+
}
363+
364+
if (this._queue.length > 0) {
365+
return this._queue.shift()!;
366+
}
367+
368+
return new Promise((resolve, reject) => {
369+
this._waiting.push({ resolve, reject });
370+
});
371+
}
372+
373+
async close(): Promise<void> {
374+
if (this._closed) return;
375+
376+
this._closed = true;
377+
this._ws.close();
378+
379+
await this._handleDisconnect();
380+
}
381+
}

0 commit comments

Comments
 (0)