Skip to content

Commit 2e5e496

Browse files
committed
fix: allow multiple subscriptions to the same URI in a single session
1 parent 5193f3a commit 2e5e496

2 files changed

Lines changed: 31 additions & 8 deletions

File tree

lib/session.ts

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ export class Session {
5757
reject: (reason: ApplicationError) => void
5858
}> = new Map();
5959
private _subscribeRequests: Map<number, SubscribeRequest> = new Map();
60-
private _subscriptions: Map<number, (event: Event) => void> = new Map();
60+
private _subscriptions: Map<number, Map<Subscription, Subscription>> = new Map();
6161
private _unsubscribeRequests: Map<number, UnsubscribeRequest> = new Map();
6262

6363
private _goodbyeRequest = (() => {
@@ -166,14 +166,24 @@ export class Session {
166166
} else if (message instanceof Subscribed) {
167167
const request = this._subscribeRequests.get(message.requestID);
168168
if (request) {
169-
this._subscriptions.set(message.subscriptionID, request.endpoint);
170-
request.promise.resolve(new Subscription(message.subscriptionID, this));
169+
const sub = new Subscription(message.subscriptionID, this, request.endpoint);
170+
let subscriptions = this._subscriptions.get(message.subscriptionID);
171+
if (!subscriptions) {
172+
subscriptions = new Map<Subscription, Subscription>();
173+
this._subscriptions.set(message.subscriptionID, subscriptions);
174+
}
175+
subscriptions.set(sub, sub);
176+
177+
request.promise.resolve(sub);
171178
this._subscribeRequests.delete(message.requestID);
172179
}
173180
} else if (message instanceof EventMsg) {
174-
const endpoint = this._subscriptions.get(message.subscriptionID);
175-
if (endpoint) {
176-
endpoint(new Event(message.args, message.kwargs, message.details));
181+
const subscriptions = this._subscriptions.get(message.subscriptionID);
182+
const event = new Event(message.args, message.kwargs, message.details);
183+
if (subscriptions) {
184+
for (const subscription of subscriptions.keys()) {
185+
subscription.eventHandler(event);
186+
}
177187
}
178188
} else if (message instanceof Unsubscribed) {
179189
const request = this._unsubscribeRequests.get(message.requestID);
@@ -365,6 +375,15 @@ export class Session {
365375
}
366376

367377
async unsubscribe(sub: Subscription): Promise<void> {
378+
const subscriptions = this._subscriptions.get(sub.subscriptionID);
379+
if (subscriptions) {
380+
subscriptions.delete(sub);
381+
if (subscriptions.size !== 0) {
382+
this._subscriptions.set(sub.subscriptionID, subscriptions);
383+
return
384+
}
385+
}
386+
368387
const unsubscribe = new Unsubscribe(new UnsubscribeFields(this._nextID, sub.subscriptionID));
369388
let promiseHandler: {
370389
resolve: () => void;

lib/types.ts

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -196,10 +196,14 @@ export class UnregisterRequest {
196196
}
197197

198198
export class Subscription {
199-
constructor(public readonly subscriptionID: number, private readonly session: Session) {}
199+
constructor(
200+
public readonly subscriptionID: number,
201+
private readonly session: Session,
202+
public readonly eventHandler: (event: Event) => void
203+
) {}
200204

201205
async unsubscribe(): Promise<void> {
202-
return this.session.unsubscribe(this)
206+
return this.session.unsubscribe(this);
203207
}
204208
}
205209

0 commit comments

Comments
 (0)