forked from mvayngrib/socketgd
-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathsocketgd.js
More file actions
318 lines (280 loc) · 8.3 KB
/
Copy pathsocketgd.js
File metadata and controls
318 lines (280 loc) · 8.3 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
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
(function(exporter) {
/**
* socket.io guaranteed delivery socket wrapper.
* if the socket gets disconnected at any point, it's up to the application to set a new socket to continue
* handling messages.
* calling 'setSocket' causes all messages that have not received an ack to be sent again.
* @constructor
*/
function SocketGD(socket, lastAcked) {
this._pending = [];
this._events = {};
this._id = 0;
this._enabled = true;
this._onAckCB = SocketGD.prototype._onAck.bind(this);
this._onReconnectCB = SocketGD.prototype._onReconnect.bind(this);
this.setLastAcked(lastAcked);
this.setSocket(socket);
}
/**
* set the id
* @param id
*/
SocketGD.prototype.setId = function (id) {
this._id = id >= 0 ? id : 0;
};
/**
* return the id
*/
SocketGD.prototype.id = function () {
return this._id;
};
/**
* set the last message id that an ack was sent for
* @param lastAcked
*/
SocketGD.prototype.setLastAcked = function(lastAcked) {
this._lastAcked = lastAcked >= 0 ? lastAcked : -1;
};
/**
* get the last acked message id
*/
SocketGD.prototype.lastAcked = function() {
return this._lastAcked;
};
/**
* replace the underlying socket.io socket with a new socket. useful in case of a socket getting
* disconnected and a new socket is used to continue with the communications
* @param socket
*/
SocketGD.prototype.setSocket = function(socket) {
this._cleanup();
this._socket = socket;
if (this._socket) {
this._socket.on('reconnect', this._onReconnectCB);
this._socket.on('socketgd_ack', this._onAckCB);
this.sendPending();
}
};
/**
* send all pending messages that have not received an ack
*/
SocketGD.prototype.sendPending = function() {
var _this = this;
// send all pending messages that haven't been acked yet
this._pending.forEach(function(message) {
_this._sendOnSocket(message);
});
};
/**
* clear out any pending messages
*/
SocketGD.prototype.clearPending = function() {
this._pending = [];
};
/**
* enable or disable sending message with gd. if disabled, then messages will be sent without guaranteeing delivery
* in case of socket disconnection/reconnection.
*/
SocketGD.prototype.enable = function(enabled) {
this._enabled = enabled;
};
/**
* get the underlying socket
*/
SocketGD.prototype.socket = function() {
return this._socket;
};
/**
* cleanup socket stuff
* @private
*/
SocketGD.prototype._cleanup = function() {
if (!this._socket) {
return;
}
this._socket.removeListener('reconnect', this._onReconnectCB);
this._socket.removeListener('socketgd_ack', this._onAckCB);
};
/**
* invoked when an ack arrives
* @param ack
* @private
*/
SocketGD.prototype._onAck = function(ack) {
// got an ack for a message, remove all messages pending an ack up to (and including) the acked message.
while (this._pending.length > 0 && this._pending[0].id <= ack.id) {
if (this._pending[0].id === ack.id && this._pending[0].ack) {
this._pending[0].ack.call(null, ack.data);
}
this._pending.shift();
}
};
/**
* invoked when an a reconnect event occurs on the underlying socket
* @private
*/
SocketGD.prototype._onReconnect = function() {
this.sendPending();
};
/**
* send an ack for a message
* @private
*/
SocketGD.prototype._sendAck = function(id, data) {
if (!this._socket) {
return;
}
this._lastAcked = id;
this._socket.emit('socketgd_ack', {id: id, data: data});
return this._lastAcked;
};
/**
* send a message on the underlying socket.io socket
* @param message
* @private
*/
SocketGD.prototype._sendOnSocket = function(message) {
if (this._enabled && message.id === undefined) {
message.id = this._id++;
message.gd = true;
this._pending.push(message);
}
if (!this._socket) {
return;
}
if (this._enabled) {
switch (message.type) {
case 'send':
this._socket.send('socketgd:' + message.id + ':' + message.msg);
break;
case 'emit':
this._socket.emit(message.event, {socketgd: message.id, msg: message.msg});
break;
}
} else {
switch (message.type) {
case 'send':
this._socket.send(message.msg, message.ack);
break;
case 'emit':
this._socket.emit(message.event, message.msg, message.ack);
break;
}
}
};
/**
* send a message with gd. this means that if an ack is not received and a new connection is established (by
* calling setSocket), the message will be sent again.
* @param message
* @param ack
*/
SocketGD.prototype.send = function(message, ack) {
this._sendOnSocket({type: 'send', msg: message, ack: ack});
};
/**
* emit an event with gd. this means that if an ack is not received and a new connection is established (by
* calling setSocket), the event will be emitted again.
* @param event
* @param message
* @param ack
*/
SocketGD.prototype.emit = function(event, message, ack) {
this._sendOnSocket({type: 'emit', event: event, msg: message, ack: ack});
};
/**
* disconnect the socket
*/
SocketGD.prototype.disconnect = function(close) {
this._socket && this._socket.disconnect(close);
this._cleanup();
this._socket = null;
};
/**
* disconnectSync the socket
*/
SocketGD.prototype.disconnectSync = function() {
this._socket && this._socket.disconnectSync();
this._cleanup();
this._socket = null;
};
/**
* close the socket
*/
SocketGD.prototype.close = function() {
this._socket && this._socket.disconnect(true);
this._cleanup();
this._socket = null;
};
SocketGD.prototype.to = function() {
this._socket = this._socket.to(...arguments);
return this;
};
/**
* listen for events on the socket. this replaces calling the 'on' method directly on the socket.io socket.
* here we take care of acking messages.
* @param event
* @param cb
*/
SocketGD.prototype.on = function(event, cb) {
this._events[event] = this._events[event] || [];
var _this = this;
var cbData = {
cb: cb,
wrapped: function(data, ack) {
if (data && event === 'message') {
// parse the message
if (data.indexOf('socketgd:') !== 0) {
cb(data, ack);
return;
}
// get the id (skipping the socketgd prefix)
var index = data.indexOf(':', 9);
if (index === -1) {
cb(data, ack);
return;
}
var id = parseInt(data.substring(9, index));
if (id <= _this._lastAcked) {
// discard the message since it was already handled and acked
return;
}
var message = data.substring(index + 1);
// the callback must call the 'ack' function so we can send an ack for the message
cb && cb(message, function(ackData) {
return _this._sendAck(id, ackData);
}, id);
} else if (data && typeof data === 'object' && data.socketgd !== undefined) {
if (data.socketgd <= _this._lastAcked) {
// discard the message since it was already handled and acked
return;
}
cb && cb(data.msg, function(ackData) {
return _this._sendAck(data.socketgd, ackData);
}, data.socketgd);
} else {
cb(data, ack);
}
}
};
this._events[event].push(cbData);
this._socket.on(event, cbData.wrapped);
};
/**
* remove a previously set callback for the specified event
*/
SocketGD.prototype.off =
SocketGD.prototype.removeListener = function(event, cb) {
if (!this._events[event]) {
return;
}
// find the callback to remove
for (var i = 0; i < this._events[event].length; ++i) {
if (this._events[event][i].cb === cb) {
this._socket && this._socket.removeListener(event, this._events[event][i].wrapped);
this._events[event].splice(i, 1);
}
}
};
exporter.SocketGD = SocketGD;
})(typeof module !== 'undefined' && typeof module.exports === 'object' ? module.exports : window);