Browse Source

Clear pending updates on socket close and process them on reconnect. Add `syncCommandId` method to align subscriptions with existing commands.

pull/15363/head
ababak 6 months ago
parent
commit
6b475fbd01
  1. 15
      ui-ngx/src/app/core/ws/telemetry-websocket.service.ts
  2. 12
      ui-ngx/src/app/core/ws/websocket.service.ts

15
ui-ngx/src/app/core/ws/telemetry-websocket.service.ts

@ -50,6 +50,7 @@ import {
UnreadCountSubCmd, UnreadCountSubCmd,
UnreadSubCmd, UnreadSubCmd,
UnsubscribeCmd, UnsubscribeCmd,
WebsocketCmd,
WebsocketDataMsg WebsocketDataMsg
} from '@app/shared/models/telemetry/telemetry.models'; } from '@app/shared/models/telemetry/telemetry.models';
import { Store } from '@ngrx/store'; import { Store } from '@ngrx/store';
@ -94,11 +95,14 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
subscriber.subscriptionCommands.forEach( subscriber.subscriptionCommands.forEach(
(subscriptionCommand) => { (subscriptionCommand) => {
if (subscriptionCommand.cmdId && (subscriptionCommand instanceof EntityDataCmd || subscriptionCommand instanceof UnreadSubCmd)) { if (subscriptionCommand.cmdId && (subscriptionCommand instanceof EntityDataCmd || subscriptionCommand instanceof UnreadSubCmd)) {
this.syncCommandId(subscriptionCommand, subscriber);
this.cmdWrapper.cmds.push(subscriptionCommand); this.cmdWrapper.cmds.push(subscriptionCommand);
} }
} }
); );
this.publishCommands(); this.publishCommands();
} else {
this.pendingUpdates.add(subscriber);
} }
} }
@ -177,4 +181,15 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
} }
} }
private syncCommandId(subscriptionCommand: WebsocketCmd, subscriber: TelemetrySubscriber) {
if (!this.subscribersMap.has(subscriptionCommand.cmdId)) {
for (const [cmdId, sub] of this.subscribersMap.entries()) {
if (sub === subscriber) {
subscriptionCommand.cmdId = cmdId;
break;
}
}
}
}
} }

12
ui-ngx/src/app/core/ws/websocket.service.ts

@ -52,6 +52,7 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
subscribersMap = new Map<number, TelemetrySubscriber | NotificationSubscriber>(); subscribersMap = new Map<number, TelemetrySubscriber | NotificationSubscriber>();
reconnectSubscribers = new Set<WsSubscriber>(); reconnectSubscribers = new Set<WsSubscriber>();
pendingUpdates = new Set<T>();
wsUri: string; wsUri: string;
@ -140,6 +141,7 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
if (close) { if (close) {
this.reconnectAttempts = 0; this.reconnectAttempts = 0;
this.lastShownCloseCode = null; this.lastShownCloseCode = null;
this.pendingUpdates.clear();
this.closeSocket(); this.closeSocket();
} }
} }
@ -223,6 +225,7 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
} }
); );
this.reconnectSubscribers.clear(); this.reconnectSubscribers.clear();
this.processPendingUpdates();
} else { } else {
this.publishCommands(); this.publishCommands();
} }
@ -298,4 +301,13 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
message, type: notificationType message, type: notificationType
})); }));
} }
private processPendingUpdates() {
if (this.pendingUpdates.size > 0) {
this.pendingUpdates.forEach((subscriber) => {
this.update(subscriber);
});
this.pendingUpdates.clear();
}
}
} }

Loading…
Cancel
Save