Browse Source

fix: add subscriberCmdIds reverse index and defensive guard in processPendingUpdates

pull/15363/head
ababak 6 months ago
parent
commit
3bd5d424e7
  1. 21
      ui-ngx/src/app/core/ws/telemetry-websocket.service.ts
  2. 11
      ui-ngx/src/app/core/ws/websocket.service.ts

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

@ -81,6 +81,12 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
const cmdId = this.nextCmdId(); const cmdId = this.nextCmdId();
if (!(subscriptionCommand instanceof MarkAsReadCmd) && !(subscriptionCommand instanceof MarkAllAsReadCmd)) { if (!(subscriptionCommand instanceof MarkAsReadCmd) && !(subscriptionCommand instanceof MarkAllAsReadCmd)) {
this.subscribersMap.set(cmdId, subscriber); this.subscribersMap.set(cmdId, subscriber);
let cmdIds = this.subscriberCmdIds.get(subscriber);
if (!cmdIds) {
cmdIds = new Set<number>();
this.subscriberCmdIds.set(subscriber, cmdIds);
}
cmdIds.add(cmdId);
} }
subscriptionCommand.cmdId = cmdId; subscriptionCommand.cmdId = cmdId;
this.cmdWrapper.cmds.push(subscriptionCommand); this.cmdWrapper.cmds.push(subscriptionCommand);
@ -141,6 +147,13 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
const cmdId = subscriptionCommand.cmdId; const cmdId = subscriptionCommand.cmdId;
if (cmdId) { if (cmdId) {
this.subscribersMap.delete(cmdId); this.subscribersMap.delete(cmdId);
const cmdIds = this.subscriberCmdIds.get(subscriber);
if (cmdIds) {
cmdIds.delete(cmdId);
if (cmdIds.size === 0) {
this.subscriberCmdIds.delete(subscriber);
}
}
} }
} }
); );
@ -183,11 +196,9 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
private syncCommandId(subscriptionCommand: WebsocketCmd, subscriber: TelemetrySubscriber) { private syncCommandId(subscriptionCommand: WebsocketCmd, subscriber: TelemetrySubscriber) {
if (!this.subscribersMap.has(subscriptionCommand.cmdId)) { if (!this.subscribersMap.has(subscriptionCommand.cmdId)) {
for (const [cmdId, sub] of this.subscribersMap.entries()) { const cmdIds = this.subscriberCmdIds.get(subscriber);
if (sub === subscriber) { if (cmdIds?.size) {
subscriptionCommand.cmdId = cmdId; subscriptionCommand.cmdId = cmdIds.values().next().value;
break;
}
} }
} }
} }

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

@ -50,6 +50,7 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
lastCmdId = 0; lastCmdId = 0;
subscribersCount = 0; subscribersCount = 0;
subscribersMap = new Map<number, TelemetrySubscriber | NotificationSubscriber>(); subscribersMap = new Map<number, TelemetrySubscriber | NotificationSubscriber>();
subscriberCmdIds = new Map<WsSubscriber, Set<number>>();
reconnectSubscribers = new Set<WsSubscriber>(); reconnectSubscribers = new Set<WsSubscriber>();
pendingUpdates = new Set<T>(); pendingUpdates = new Set<T>();
@ -136,6 +137,7 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
} }
this.lastCmdId = 0; this.lastCmdId = 0;
this.subscribersMap.clear(); this.subscribersMap.clear();
this.subscriberCmdIds.clear();
this.subscribersCount = 0; this.subscribersCount = 0;
this.cmdWrapper.clear(); this.cmdWrapper.clear();
if (close) { if (close) {
@ -304,10 +306,13 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
private processPendingUpdates() { private processPendingUpdates() {
if (this.pendingUpdates.size > 0) { if (this.pendingUpdates.size > 0) {
this.pendingUpdates.forEach((subscriber) => { const subscribers = Array.from(this.pendingUpdates);
this.update(subscriber);
});
this.pendingUpdates.clear(); this.pendingUpdates.clear();
for (const subscriber of subscribers) {
if (this.subscriberCmdIds.has(subscriber)) {
this.update(subscriber);
}
}
} }
} }
} }

Loading…
Cancel
Save