Skip to content

Commit 109ad35

Browse files
authored
Merge pull request #34 from xconnio/connected-function-disconnect-callback
Add function to check connection state, handle goodbye and disconnect callback
2 parents 701f807 + 47622c0 commit 109ad35

4 files changed

Lines changed: 130 additions & 3 deletions

File tree

examples/reconnect.ts

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,55 @@
1+
import { connectAnonymous } from "../lib";
2+
3+
async function sleep(ms: number): Promise<void> {
4+
return new Promise(resolve => setTimeout(resolve, ms));
5+
}
6+
7+
async function main() {
8+
let retryDelay = 1000; // 1 second
9+
const maxDelay = 30000; // 30 seconds
10+
11+
for (; ;) {
12+
try {
13+
const session = await connectAnonymous("ws://localhost:8080/ws", "realm1");
14+
console.log("Connected successfully");
15+
16+
// Reset backoff after successful connection
17+
retryDelay = 1000;
18+
19+
// Register multiple disconnect callbacks
20+
session.onDisconnect(async () => {
21+
console.log("Callback 1: disconnection event.");
22+
});
23+
24+
session.onDisconnect(async () => {
25+
console.log("Callback 2: disconnection event.");
26+
await sleep(500);
27+
});
28+
29+
session.onDisconnect(async () => {
30+
console.log("Callback 3: disconnection event.");
31+
});
32+
33+
// Wait for disconnect
34+
const disconnected = new Promise<void>(resolve => {
35+
session.onDisconnect(async () => {
36+
console.log("Disconnected from router!");
37+
resolve();
38+
});
39+
});
40+
41+
await disconnected;
42+
console.log("Retrying connection...");
43+
44+
} catch (err) {
45+
console.error(`Failed to connect: ${err}`);
46+
console.log(`Retrying in ${retryDelay / 1000}s...`);
47+
await sleep(retryDelay);
48+
49+
// Exponential backoff
50+
retryDelay = Math.min(retryDelay * 2, maxDelay);
51+
}
52+
}
53+
}
54+
55+
main().catch(err => console.error(`Fatal error: ${err}`));

lib/session.ts

Lines changed: 54 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,11 +18,12 @@ import {
1818
Event as EventMsg,
1919
Unsubscribe, UnsubscribeFields,
2020
Unsubscribed,
21-
Error, ErrorFields
21+
Error, ErrorFields,
22+
Goodbye, GoodbyeFields
2223
} from "wampproto";
2324

2425
import {wampErrorString} from "./helpers";
25-
import {ERROR_RUNTIME_ERROR} from "./wamp";
26+
import {ERROR_RUNTIME_ERROR, CLOSE_CLOSE_REALM} from "./wamp";
2627
import {ApplicationError, ProtocolError} from "./exception";
2728
import {
2829
IBaseSession,
@@ -42,6 +43,7 @@ export class Session {
4243
private _baseSession: IBaseSession;
4344
private _wampSession: WAMPSession;
4445
private _idGen: SessionScopeIDGenerator = new SessionScopeIDGenerator();
46+
private _disconnectCallbacks: Array<() => Promise<void>> = [];
4547

4648
private _callRequests: Map<number, {
4749
resolve: (value: Result) => void,
@@ -58,9 +60,24 @@ export class Session {
5860
private _subscriptions: Map<number, (event: Event) => void> = new Map();
5961
private _unsubscribeRequests: Map<number, UnsubscribeRequest> = new Map();
6062

63+
private _goodbyeRequest = (() => {
64+
let resolve!: () => void;
65+
let isCompleted = false;
66+
const promise = new Promise<void>((res) => {
67+
resolve = () => {
68+
if (!isCompleted) {
69+
isCompleted = true;
70+
res();
71+
}
72+
};
73+
});
74+
return { promise, resolve, isCompleted };
75+
})();
76+
6177
constructor(baseSession: IBaseSession) {
6278
this._baseSession = baseSession;
6379
this._wampSession = new WAMPSession(baseSession.serializer());
80+
this._baseSession.onDisconnect(async () => { await this.markDisconnected();});
6481

6582
(async () => {
6683
for (; ;) {
@@ -70,12 +87,34 @@ export class Session {
7087
})();
7188
}
7289

90+
onDisconnect(callback: () => Promise<void>): void {
91+
this._disconnectCallbacks.push(callback);
92+
}
93+
7394
private get _nextID(): number {
7495
return this._idGen.next();
7596
}
7697

7798
async close(): Promise<void> {
78-
await this._baseSession.close();
99+
const goodbye = new Goodbye(new GoodbyeFields({}, CLOSE_CLOSE_REALM));
100+
const data = this._wampSession.sendMessage(goodbye)
101+
this._baseSession.send(data)
102+
103+
return Promise.race([
104+
this._goodbyeRequest.promise,
105+
new Promise<void>((resolve) =>
106+
setTimeout(async () => {
107+
await this._baseSession.close();
108+
resolve();
109+
}, 10_000)
110+
)
111+
]).finally(async () => {
112+
await this._baseSession.close();
113+
});
114+
}
115+
116+
isConnected(): boolean {
117+
return this._baseSession.isConnected();
79118
}
80119

81120
private async _processIncomingMessage(message: Message): Promise<void> {
@@ -200,11 +239,23 @@ export class Session {
200239
default:
201240
throw new ProtocolError(wampErrorString(message));
202241
}
242+
} else if (message instanceof Goodbye) {
243+
await this.markDisconnected()
203244
} else {
204245
throw new ProtocolError(`Unexpected message type ${typeof message}`);
205246
}
206247
}
207248

249+
private async markDisconnected() {
250+
if (this._disconnectCallbacks.length > 0) {
251+
await Promise.all(this._disconnectCallbacks.map(cb => cb()));
252+
}
253+
254+
if (this._goodbyeRequest && !this._goodbyeRequest.isCompleted) {
255+
this._goodbyeRequest.resolve();
256+
}
257+
}
258+
208259
async call(
209260
procedure: string,
210261
args?: any[] | null,

lib/types.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,13 +44,22 @@ export abstract class IBaseSession {
4444
async close(): Promise<void> {
4545
throw new Error("UnimplementedError");
4646
}
47+
48+
isConnected(): boolean {
49+
throw new Error("UnimplementedError");
50+
}
51+
52+
onDisconnect(callback: () => Promise<void>): void {
53+
throw new Error("UnimplementedError");
54+
}
4755
}
4856

4957
export class BaseSession extends IBaseSession {
5058
private readonly _ws: WebSocket;
5159
private readonly _wsMessageHandler: any;
5260
private readonly sessionDetails: SessionDetails;
5361
private readonly _serializer: Serializer;
62+
private _disconnectCallbacks: Array<() => Promise<void>> = [];
5463

5564
constructor(
5665
ws: WebSocket,
@@ -66,6 +75,9 @@ export class BaseSession extends IBaseSession {
6675

6776
// close cleanly on abrupt client disconnect
6877
this._ws.addEventListener("close", async () => {
78+
if (this._disconnectCallbacks.length > 0) {
79+
await Promise.all(this._disconnectCallbacks.map(cb => cb()));
80+
}
6981
await this.close();
7082
});
7183
}
@@ -126,6 +138,14 @@ export class BaseSession extends IBaseSession {
126138

127139
this._ws.close();
128140
}
141+
142+
isConnected(): boolean {
143+
return this._ws.readyState === WebSocket.OPEN;
144+
}
145+
146+
onDisconnect(callback: () => Promise<void>): void {
147+
this._disconnectCallbacks.push(callback);
148+
}
129149
}
130150

131151
export class Result {

lib/wamp.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
1+
export const CLOSE_CLOSE_REALM = "wamp.close.close_realm"
12
export const ERROR_RUNTIME_ERROR = "wamp.error.runtime_error"

0 commit comments

Comments
 (0)