More work on desktop multiplexor.

Ylian Saint-Hilaire committed Apr 27, 2020 at 16:25 UTC 11ccefc61295cd3f68bf05784e4dc9acb1e992bf
2 files changed +123 -52
meshdesktopmultiplex.js
+122 -52
@@ -86,13 +86,11 @@ function CreateDesktopDecoder() {
86 // Add an agent or viewer
87 obj.addPeer = function (peer) {
88 if (peer.req.query.browser) {
89 - console.log('addPeer-viewer');
90 -
91 - // This is a viewer
89 + //console.log('addPeer-viewer');
90 +
91 + // Setup the viewer
92 if (obj.viewers.indexOf(peer) >= 0) return true;
93 obj.viewers.push(peer);
94 -
95 - // Setup the viewer
94 peer.desktopPaused = true;
95 peer.imageCompression = 30;
96 peer.imageScaling = 1024;
@@ -101,19 +99,22 @@ function CreateDesktopDecoder() {
99 peer.dataPtr = obj.firstData;
100 peer.sending = false;
101 peer.sendQueue = [];
102 + peer.paused = false;
103 +
104 + // Indicated we are connected
105 + obj.sendToViewer(peer, 'c');
106 } else {
105 - console.log('addPeer-agent');
107 + //console.log('addPeer-agent');
108 if (obj.agent != null) return false;
109
108 - // This is the agent
109 - obj.agent = peer;
110 -
110 // Setup the agent
111 + obj.agent = peer;
112 peer.sending = false;
113 peer.sendQueue = [];
114 - //peer.ws.send('{"tsid":10,"type":"options"}');
115 - //peer.ws.send('2');
114 + peer.paused = false;
115
116 + // Indicated we are connected and send connection options and protocol if needed
117 + obj.sendToAgent('c');
118 if (obj.viewerConnected == true) {
119 if (obj.protocolOptions != null) { obj.sendToAgent(JSON.stringify(obj.protocolOptions)); } // Send connection options
120 obj.sendToAgent('2'); // Send remote desktop connect
@@ -123,23 +124,52 @@ function CreateDesktopDecoder() {
124 }
125
126 // Remove an agent or viewer
127 + // Return true if this multiplexor is no longer needed.
128 obj.removePeer = function (peer) {
127 - if (peer == agent) {
128 - console.log('removePeer-agent');
129 + if (peer == obj.agent) {
130 + //console.log('removePeer-agent');
131 // Clean up the agent
132 + obj.agent = null;
133
134 // Agent has disconnected, disconnect everyone.
135 + for (var i in obj.viewers) { obj.viewers[i].close(); }
136 + dispose();
137 + return true;
138 } else {
133 - console.log('removePeer-viewer');
139 + //console.log('removePeer-viewer');
140 // Remove a viewer
141 var i = obj.viewers.indexOf(peer);
142 if (i == -1) return false;
143 obj.viewers.splice(i, 1);
144
139 - // Clean up the viewer
140 -
145 + // Aggressive clean up of the viewer
146 + delete peer.desktopPaused;
147 + delete peer.imageCompression;
148 + delete peer.imageScaling;
149 + delete peer.imageFrameRate;
150 + delete peer.lastImageNumberSent;
151 + delete peer.dataPtr;
152 + delete peer.sending;
153 + delete peer.sendQueue;
154 +
155 + // Resume flow control if this was the peer
156 + if (peer.sending == true) {
157 + obj.viewersSendingCount--;
158 + peer.sending = false;
159 + if ((obj.viewersSendingCount < obj.viewers.length) && obj.agent && (obj.agent.paused == true)) { obj.agent.paused = false; obj.agent.ws._socket.resume(); }
160 + }
161 +
162 + // If this is the last viewer, disconnect the agent
163 + if ((obj.viewers.length == 0) && (obj.agent != null)) { obj.agent.close(); dispose(); return true; }
164 }
142 - return true;
165 + return false;
166 + }
167 +
168 + // Clean up ourselves
169 + function dispose() {
170 + delete obj.viewers;
171 + delete obj.imagesCounters;
172 + delete obj.images;
173 }
174
175 // Send data to the agent or queue it up for sending
@@ -148,21 +178,32 @@ function CreateDesktopDecoder() {
178 //console.log('SendToAgent', data.length);
179 if (obj.agent.sending) {
180 obj.agent.sendQueue.push(data);
151 - // TODO: Flow control, stop all viewers
181 } else {
182 obj.agent.ws.send(data, sendAgentNext);
183 +
184 + // Flow control, pause all viewers
185 + for (var i in obj.viewers) {
186 + var v = obj.viewers[i];
187 + if (v.paused == false) { v.paused = true; v.ws._socket.pause(); }
188 + }
189 }
190 }
191
192 // Send more data to the agent
193 function sendAgentNext() {
194 + if (obj.agent == null) return;
195 if (obj.agent.sendQueue.length > 0) {
196 // Send from the pending send queue
197 obj.agent.ws.send(obj.agent.sendQueue.shift(), sendAgentNext);
198 } else {
199 // Nothing to send
200 obj.agent.sending = false;
165 - // TODO: Flow control, start all viewers
201 +
202 + // Flow control, resume all viewers
203 + for (var i in obj.viewers) {
204 + var v = obj.viewers[i];
205 + if (v.paused == true) { v.paused = false; v.ws._socket.resume(); }
206 + }
207 }
208 }
209
@@ -177,16 +218,19 @@ function CreateDesktopDecoder() {
218 //console.log('SendToViewer', data.length);
219 if (viewer.sending) {
220 viewer.sendQueue.push(data);
180 - // TODO: Flow control, stop the agent
221 } else {
222 viewer.sending = true;
183 - obj.viewersSendingCount++;
223 viewer.ws.send(data, function () { sendViewerNext(viewer); });
224 +
225 + // Flow control, pause the agent if needed
226 + obj.viewersSendingCount++;
227 + if ((obj.viewersSendingCount >= obj.viewers.length) && obj.agent && (obj.agent.paused == false)) { obj.agent.paused = true; obj.agent.ws._socket.pause(); }
228 }
229 }
230
231 // Send more data to the viewer
232 function sendViewerNext(viewer) {
233 + if (viewer.sendQueue == null) return;
234 if (viewer.sendQueue.length > 0) {
235 // Send from the pending send queue
236 if (viewer.sending == false) { viewer.sending = true; obj.viewersSendingCount++; }
@@ -199,13 +243,21 @@ function CreateDesktopDecoder() {
243 viewer.lastImageNumberSent = viewer.dataPtr;
244 if ((image.next != null) && ((viewer.dataPtr + 1) != image.next)) { console.log('SVIEW-S2', viewer.dataPtr, image.next); } // DEBUG
245 viewer.dataPtr = image.next;
202 - if (viewer.sending == false) { viewer.sending = true; obj.viewersSendingCount++; }
246 viewer.ws.send(image.data, function () { sendViewerNext(viewer); });
247 +
248 + // Flow control, pause the agent if needed
249 + if (viewer.sending == false) {
250 + viewer.sending = true;
251 + obj.viewersSendingCount++;
252 + if ((obj.viewersSendingCount >= obj.viewers.length) && obj.agent && (obj.agent.paused == false)) { obj.agent.paused = true; obj.agent.ws._socket.pause(); }
253 + }
254 } else {
255 // Nothing to send
256 viewer.sending = false;
257 +
258 + // Flow control, resume agent if needed
259 obj.viewersSendingCount--;
208 - // TODO: Flow control, start agent
260 + if ((obj.viewersSendingCount < obj.viewers.length) && obj.agent && (obj.agent.paused == true)) { obj.agent.paused = false; obj.agent.ws._socket.resume(); }
261 }
262 }
263 }
@@ -470,7 +522,7 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
522 }
523
524 // If there is no authentication, drop this connection
473 - if ((obj.id != null) && (obj.id.startsWith('meshmessenger/') == false) && (obj.user == null) && (obj.ruserid == null)) { try { ws.close(); parent.parent.debug('relay', 'Relay: Connection with no authentication (' + cleanRemoteAddr(obj.req.ip) + ')'); } catch (e) { console.log(e); } return; }
525 + if ((obj.id != null) && (obj.user == null) && (obj.ruserid == null)) { try { ws.close(); parent.parent.debug('relay', 'Relay: Connection with no authentication (' + cleanRemoteAddr(obj.req.ip) + ')'); } catch (e) { console.log(e); } return; }
526
527 // Relay session count (we may remove this in the future)
528 obj.relaySessionCounted = true;
@@ -500,13 +552,25 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
552
553 // Disconnect this agent
554 obj.close = function (arg) {
555 + if (obj.ws == null) return; // Already closed.
556 +
557 + // Close the connection
558 if ((arg == 1) || (arg == null)) { try { ws.close(); parent.parent.debug('relay', 'Relay: Soft disconnect (' + cleanRemoteAddr(obj.req.ip) + ')'); } catch (e) { console.log(e); } } // Soft close, close the websocket
559 if (arg == 2) { try { ws._socket._parent.end(); parent.parent.debug('relay', 'Relay: Hard disconnect (' + cleanRemoteAddr(obj.req.ip) + ')'); } catch (e) { console.log(e); } } // Hard close, close the TCP socket
560 + if (obj.relaySessionCounted) { parent.relaySessionCount--; delete obj.relaySessionCounted; }
561 + if (obj.deskDecoder != null) { if (obj.deskDecoder.removePeer(obj) == true) { delete parent.desktoprelays[obj.id]; } }
562
563 // Aggressive cleanup
564 delete obj.id;
565 delete obj.ws;
509 - delete obj.peer;
566 + delete obj.req;
567 + delete obj.user;
568 + delete obj.ruserid;
569 + delete obj.deskDecoder;
570 +
571 + // Clear timers if present
572 + if (obj.pingtimer != null) { clearInterval(obj.pingtimer); delete obj.pingtimer; }
573 + if (obj.pongtimer != null) { clearInterval(obj.pongtimer); delete obj.pongtimer; }
574 };
575
576 obj.sendAgentMessage = function (command, userid, domainid) {
@@ -572,6 +636,31 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
636 if (obj.id == null) { try { obj.close(); } catch (e) { } return null; } // Attempt to connect without id, drop this.
637 ws._socket.setKeepAlive(true, 240000); // Set TCP keep alive
638
639 + // Validate that the id is valid, we only need to do this on non-authenticated sessions.
640 + // TODO: Figure out when this needs to be done.
641 + if ((user == null) && (!parent.args.notls)) {
642 + // Check the identifier, if running without TLS, skip this.
643 + var ids = obj.id.split(':');
644 + if (ids.length != 3) { ws.close(); delete obj.id; return null; } // Invalid ID, drop this.
645 + if (parent.crypto.createHmac('SHA384', parent.relayRandom).update(ids[0] + ':' + ids[1]).digest('hex') != ids[2]) { ws.close(); delete obj.id; return null; } // Invalid HMAC, drop this.
646 + if ((Date.now() - parseInt(ids[1])) > 120000) { ws.close(); delete obj.id; return null; } // Expired time, drop this.
647 + obj.id = ids[0];
648 + }
649 +
650 + // Setup the agent PING/PONG timers
651 + if ((typeof parent.parent.args.agentping == 'number') && (obj.pingtimer == null)) { obj.pingtimer = setInterval(sendPing, parent.parent.args.agentping * 1000); }
652 + else if ((typeof parent.parent.args.agentpong == 'number') && (obj.pongtimer == null)) { obj.pongtimer = setInterval(sendPong, parent.parent.args.agentpong * 1000); }
653 +
654 + // Create if needed and add this peer to the desktop multiplexor
655 + obj.deskDecoder = parent.desktoprelays[obj.id];
656 + if (obj.deskDecoder == null) {
657 + obj.deskDecoder = CreateDesktopDecoder();
658 + parent.desktoprelays[obj.id] = obj.deskDecoder;
659 + }
660 + obj.deskDecoder.addPeer(obj);
661 + ws._socket.resume(); // Release the traffic
662 +
663 + /*
664 // If this is a MeshMessenger session, the ID is the two userid's and authentication must match one of them.
665 if (obj.id.startsWith('meshmessenger/')) {
666 if ((obj.id.startsWith('meshmessenger/user/') == true) && (user == null)) { try { obj.close(); } catch (e) { } return null; } // If user-to-user, both sides need to be authenticated.
@@ -588,7 +677,6 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
677
678 // Validate that the id is valid, we only need to do this on non-authenticated sessions.
679 // TODO: Figure out when this needs to be done.
591 - /*
680 if (!parent.args.notls) {
681 // Check the identifier, if running without TLS, skip this.
682 var ids = obj.id.split(':');
@@ -597,7 +685,6 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
685 if ((Date.now() - parseInt(ids[1])) > 120000) { ws.close(); delete obj.id; return null; } // Expired time, drop this.
686 obj.id = ids[0];
687 }
600 - */
688
689 // Check the peer connection status
690 {
@@ -746,48 +833,30 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
833 }
834 }
835 }
836 + */
837 }
838
751 - ws.flushSink = function () { try { ws._socket.resume(); } catch (ex) { console.log(ex); } };
752 -
839 // When data is received from the mesh relay web socket
840 ws.on('message', function (data) {
755 - // If this data was received by the agent, decode it.
841 + // If this data was received by the agent, decode it.
842 if (this.me.deskDecoder != null) { this.me.deskDecoder.processData(this.me, data); }
757 -
758 - /*
759 - //console.log(typeof data, data.length);
760 - if (this.peer != null) {
761 - //if (typeof data == 'string') { console.log('Relay: ' + data); } else { console.log('Relay:' + data.length + ' byte(s)'); }
762 - try {
763 - this._socket.pause();
764 - if (this.logfile != null) {
765 - // Write data to log file then perform relay
766 - var xthis = this;
767 - recordingEntry(this.logfile.fd, 2, ((obj.req.query.browser) ? 2 : 0), data, function () { xthis.peer.send(data, ws.flushSink); });
768 - } else {
769 - // Perform relay
770 - this.peer.send(data, ws.flushSink);
771 - }
772 - } catch (ex) { console.log(ex); }
773 - }
774 - */
843 });
844
845 // If error, close both sides of the relay.
846 ws.on('error', function (err) {
847 + //console.log('ws-error', err);
848 parent.relaySessionErrorCount++;
780 - if (obj.relaySessionCounted) { parent.relaySessionCount--; delete obj.relaySessionCounted; }
849 console.log('Relay error from ' + cleanRemoteAddr(obj.req.ip) + ', ' + err.toString().split('\r')[0] + '.');
782 - closeBothSides();
850 + obj.close();
851 });
852
853 // If the relay web socket is closed, close both sides.
854 ws.on('close', function (req) {
787 - if (obj.relaySessionCounted) { parent.relaySessionCount--; delete obj.relaySessionCounted; }
788 - closeBothSides();
855 + //console.log('ws-close', req);
856 + obj.close();
857 });
858
859 + /*
860 // Close both our side and the peer side.
861 function closeBothSides() {
862 if (obj.id != null) {
@@ -876,6 +945,7 @@ module.exports.CreateMeshRelay = function (parent, ws, req, domain, user, cookie
945 }
946 } catch (ex) { console.log(ex); func(fd, tag); }
947 }
948 + */
949
950 // Mark this relay session as authenticated if this is the user end.
951 obj.authenticated = (user != null);
webserver.js
+1
@@ -172,6 +172,7 @@ module.exports.CreateWebServer = function (parent, db, args, certificates) {
172 obj.wsPeerSessions3 = {}; // ServerId --> UserId --> [ SessionId ]
173 obj.sessionsCount = {}; // Merged session counters, used when doing server peering. UserId --> SessionCount
174 obj.wsrelays = {}; // Id -> Relay
175 + obj.desktoprelays = {}; // Id -> Desktop Multiplexor Relay
176 obj.wsPeerRelays = {}; // Id -> { ServerId, Time }
177 var tlsSessionStore = {}; // Store TLS session information for quick resume.
178 var tlsSessionStoreCount = 0; // Number of cached TLS session information in store.