diff --git a/src/socket/socketuser.js b/src/socket/socketuser.js index b63392d..d3d2550 100644 --- a/src/socket/socketuser.js +++ b/src/socket/socketuser.js @@ -271,9 +271,17 @@ export class SocketUser { async handleDisconnect() { await this.actionManager.removeAllListeners(); + await this.updateManager.removeAllListeners(); await this.eventManager.removeAllListeners(); await this.statsManager.removeAllListeners(); await this.notificationManager.removeAllListeners(); - logger.info('External user disconnected:', this.socket.user?.username); + if (this.user || this.socket.user) { + logger.info( + 'External user disconnected:', + this.user.username || this.socket.user?.username + ); + } else { + logger.info('External user disconnected.'); + } } } diff --git a/src/updates/__tests__/updatemanager.test.js b/src/updates/__tests__/updatemanager.test.js index cd93e39..5369338 100644 --- a/src/updates/__tests__/updatemanager.test.js +++ b/src/updates/__tests__/updatemanager.test.js @@ -211,4 +211,19 @@ describe('UpdateManager', () => { ); }); }); + + describe('removeAllListeners', () => { + it('should remove all subscriptions', async () => { + await updateManager.subscribeToObjectNew('printer'); + await updateManager.subscribeToObjectDelete('printer'); + await updateManager.subscribeToObjectUpdate('123', 'printer'); + + expect(updateManager.subscriptions.size).toBe(3); + + await updateManager.removeAllListeners(); + + expect(natsServer.removeSubscription).toHaveBeenCalledTimes(3); + expect(updateManager.subscriptions.size).toBe(0); + }); + }); }); diff --git a/src/updates/updatemanager.js b/src/updates/updatemanager.js index 66b904d..ee73492 100644 --- a/src/updates/updatemanager.js +++ b/src/updates/updatemanager.js @@ -78,12 +78,15 @@ const stableStringify = value => { const getSubscriptionOwner = (socketId, filter) => `${socketId}:${stableStringify(normalizeFilter(filter))}`; +const getSubscriptionKey = (subject, owner) => `${subject}:${owner}`; + /** * UpdateManager handles tracking object updates and broadcasts update events via websockets. */ export class UpdateManager { constructor(socketClient) { this.socketClient = socketClient; + this.subscriptions = new Set(); } matchesObjectTypeFilter(objectType, filter, value) { @@ -116,79 +119,108 @@ export class UpdateManager { async subscribeToObjectNew(objectType, filter = {}) { const normalizedFilter = normalizeFilter(filter); - await natsServer.subscribe( - `${objectType}s.new`, - getSubscriptionOwner(this.socketClient.socketId, normalizedFilter), - async (key, value) => { - logger.trace('Object new event:', value); - this.emitObjectTypeEvent( - 'objectNew', - objectType, - normalizedFilter, - value - ); - } + const subject = `${objectType}s.new`; + const owner = getSubscriptionOwner( + this.socketClient.socketId, + normalizedFilter ); + + await natsServer.subscribe(subject, owner, async (key, value) => { + logger.trace('Object new event:', value); + this.emitObjectTypeEvent( + 'objectNew', + objectType, + normalizedFilter, + value + ); + }); + + this.subscriptions.add(getSubscriptionKey(subject, owner)); return { success: true }; } async subscribeToObjectDelete(objectType, filter = {}) { const normalizedFilter = normalizeFilter(filter); - await natsServer.subscribe( - `${objectType}s.delete`, - getSubscriptionOwner(this.socketClient.socketId, normalizedFilter), - async (key, value) => { - logger.trace('Object delete event:', value); - this.emitObjectTypeEvent( - 'objectDelete', - objectType, - normalizedFilter, - value - ); - } + const subject = `${objectType}s.delete`; + const owner = getSubscriptionOwner( + this.socketClient.socketId, + normalizedFilter ); + + await natsServer.subscribe(subject, owner, async (key, value) => { + logger.trace('Object delete event:', value); + this.emitObjectTypeEvent( + 'objectDelete', + objectType, + normalizedFilter, + value + ); + }); + + this.subscriptions.add(getSubscriptionKey(subject, owner)); return { success: true }; } async subscribeToObjectUpdate(id, objectType) { logger.debug('Subscribing to object update...', id, objectType); - await natsServer.subscribe( - `${objectType}s.${id}.object`, - this.socketClient.socketId, - (key, value) => { - const expandedValue = expandObjectIds(value); - logger.trace('Object update event:', id); - this.socketClient.socket.emit('objectUpdate', { - _id: id, - objectType: objectType, - object: { ...expandedValue } - }); - } - ); + const subject = `${objectType}s.${id}.object`; + const owner = this.socketClient.socketId; + + await natsServer.subscribe(subject, owner, (key, value) => { + const expandedValue = expandObjectIds(value); + logger.trace('Object update event:', id); + this.socketClient.socket.emit('objectUpdate', { + _id: id, + objectType: objectType, + object: { ...expandedValue } + }); + }); + + this.subscriptions.add(getSubscriptionKey(subject, owner)); return { success: true }; } async removeObjectNewListener(objectType, filter = {}) { - await natsServer.removeSubscription( - `${objectType}s.new`, - getSubscriptionOwner(this.socketClient.socketId, filter) - ); + const subject = `${objectType}s.new`; + const owner = getSubscriptionOwner(this.socketClient.socketId, filter); + + await natsServer.removeSubscription(subject, owner); + this.subscriptions.delete(getSubscriptionKey(subject, owner)); return { success: true }; } async removeObjectDeleteListener(objectType, filter = {}) { - await natsServer.removeSubscription( - `${objectType}s.delete`, - getSubscriptionOwner(this.socketClient.socketId, filter) - ); + const subject = `${objectType}s.delete`; + const owner = getSubscriptionOwner(this.socketClient.socketId, filter); + + await natsServer.removeSubscription(subject, owner); + this.subscriptions.delete(getSubscriptionKey(subject, owner)); return { success: true }; } async removeObjectUpdateListener(id, objectType) { - await natsServer.removeSubscription( - `${objectType}s.${id}.object`, - this.socketClient.socketId + const subject = `${objectType}s.${id}.object`; + const owner = this.socketClient.socketId; + + await natsServer.removeSubscription(subject, owner); + this.subscriptions.delete(getSubscriptionKey(subject, owner)); + return { success: true }; + } + + async removeAllListeners() { + logger.debug('Removing all update listeners...'); + const removePromises = Array.from(this.subscriptions).map( + subscriptionKey => { + const separatorIndex = subscriptionKey.indexOf(':'); + const subject = subscriptionKey.slice(0, separatorIndex); + const owner = subscriptionKey.slice(separatorIndex + 1); + return natsServer.removeSubscription(subject, owner); + } ); + + await Promise.all(removePromises); + this.subscriptions.clear(); + logger.debug(`Removed ${removePromises.length} update listener(s)`); return { success: true }; } }