/
clientEntityContext.ts
364 lines (327 loc) · 12 KB
/
clientEntityContext.ts
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
// Copyright (c) Microsoft Corporation. All rights reserved.
// Licensed under the MIT License.
import * as log from "./log";
import { StreamingReceiver } from "./core/streamingReceiver";
import { MessageSender } from "./core/messageSender";
import { ManagementClient, ManagementClientOptions } from "./core/managementClient";
import { ConnectionContext } from "./connectionContext";
import { Dictionary, AmqpError } from "rhea-promise";
import { ClientType } from "./client";
import { BatchingReceiver } from "./core/batchingReceiver";
import { ConcurrentExpiringMap } from "./util/concurrentExpiringMap";
import { MessageReceiver } from "./core/messageReceiver";
import { MessageSession } from "./session/messageSession";
import { SessionManager } from "./session/sessionManager";
/**
* @interface ClientEntityContext
* Provides contextual information like the underlying amqp connection, cbs session,
* management session, tokenProvider, senders, receivers, etc. about the ServiceBus client.
* @internal
*/
export interface ClientEntityContextBase {
/**
* @property {ConnectionContext} namespace Describes the context with common properties at
* the namespace level.
*/
namespace: ConnectionContext;
/**
* @property {string} entityPath - The name/path of the entity (queue/topic/subscription) to which
* the connection needs to happen.
*/
entityPath: string;
/**
* @property {boolean} [isSessionEnabled] Indicates whether the client entity is session enabled.
* Default: `false`.
*/
isSessionEnabled?: boolean;
/**
* @property {ManagementClient} [managementClient] A reference to the management client
* ($management endpoint) on the underlying amqp connection for the ServiceBus Client.
*/
managementClient?: ManagementClient;
/**
* @property {StreamingReceiver} [receiver] The ServiceBus receiver associated with the
* client entity for streaming messages.
*/
streamingReceiver?: StreamingReceiver;
/**
* @property {BatchingReceiver} [batchingReceiver] The ServiceBus receiver associated with the
* client entity for receiving a batch of messages.
*/
batchingReceiver?: BatchingReceiver;
/**
* @property {Dictionary<MessageSession>} messageSessions A dictionary of the MessageSession
* objects associated with this client.
*/
messageSessions: Dictionary<MessageSession>;
/**
* @property {Dictionary<MessageSession>} expiredMessageSessions A dictionary that stores expired message sessions IDs.
*/
expiredMessageSessions: Dictionary<boolean>;
/**
* @property {MessageSender} [sender] The ServiceBus sender associated with the client entity.
*/
sender?: MessageSender;
/**
* @property {ConcurrentExpiringMap<string>} [requestResponseLockedMessages] A map of locked
* messages received using the management client.
*/
requestResponseLockedMessages: ConcurrentExpiringMap<string>;
/**
* @property {SessionManager} [sessionManager] SessionManager is responsible for efficiently
* receiving messages from multiple message sessions.
*/
sessionManager?: SessionManager;
/**
* @property {ClientType} [clientType] Type of the client, used mostly for logging
*/
clientType: ClientType;
/**
* @property {string} [clientId] Unique Id of the client for which this context is created
*/
clientId: string;
/**
* @property {boolean} [isClosed] Denotes if close() was called on this client.
*/
isClosed: boolean;
}
/**
* @internal
*/
export interface ClientEntityContext extends ClientEntityContextBase {
onDetached(error?: AmqpError | Error): Promise<void>;
getReceiver(name: string, sessionId?: string): MessageReceiver | MessageSession | undefined;
close(): Promise<void>;
}
/**
* @internal
*/
export interface ClientEntityContextOptions {
managementClientAddress?: string;
managementClientAudience?: string;
isSessionEnabled?: boolean;
}
/**
* @internal
*/
export namespace ClientEntityContext {
/**
* @internal
*/
export function create(
entityPath: string,
clientType: ClientType,
context: ConnectionContext,
clientId: string,
options?: ClientEntityContextOptions
): ClientEntityContext {
log.entityCtxt(
"[%s] Creating client entity context for %s: %O",
context.connectionId,
clientId
);
if (!options) options = {};
const entityContext: ClientEntityContextBase = {
namespace: context,
entityPath: entityPath,
clientType: clientType,
clientId: clientId,
isClosed: false,
requestResponseLockedMessages: new ConcurrentExpiringMap<string>(),
isSessionEnabled: !!options.isSessionEnabled,
messageSessions: {},
expiredMessageSessions: {}
};
(entityContext as ClientEntityContext).sessionManager = new SessionManager(
entityContext as ClientEntityContext
);
(entityContext as ClientEntityContext).getReceiver = (name: string, sessionId?: string) => {
if (sessionId != undefined && entityContext.expiredMessageSessions[sessionId]) {
const error = new Error(
`The session lock has expired on the session with id ${sessionId}.`
);
error.name = "SessionLockLostError";
log.error(
"[%s] Failed to find receiver '%s' as the session with id '%s' is expired",
entityContext.namespace.connectionId,
name,
sessionId
);
throw error;
}
if (
sessionId != null &&
entityContext.messageSessions[sessionId] &&
entityContext.messageSessions[sessionId].name === name
) {
return entityContext.messageSessions[sessionId];
}
if (entityContext.streamingReceiver && entityContext.streamingReceiver.name === name) {
return entityContext.streamingReceiver;
}
if (entityContext.batchingReceiver && entityContext.batchingReceiver.name === name) {
return entityContext.batchingReceiver;
}
let existingReceivers = "";
if (sessionId != null && entityContext.messageSessions[sessionId]) {
existingReceivers = entityContext.messageSessions[sessionId].name;
} else {
if (entityContext.streamingReceiver) {
existingReceivers = entityContext.streamingReceiver.name;
}
if (entityContext.batchingReceiver) {
existingReceivers +=
(existingReceivers ? ", " : "") + entityContext.batchingReceiver.name;
}
}
log.error(
"[%s] Failed to find receiver '%s' among existing receivers: %s",
entityContext.namespace.connectionId,
name,
existingReceivers
);
return;
};
(entityContext as ClientEntityContext).onDetached = async (error?: AmqpError | Error) => {
const connectionId = entityContext.namespace.connectionId;
const detachCalls: Promise<void>[] = [];
// Call onDetached() on sender so that it can decide whether to reconnect or not
const sender = entityContext.sender;
if (sender && !sender.isConnecting) {
log.error("[%s] calling detached on sender '%s'.", connectionId, sender.name);
detachCalls.push(
sender.onDetached().catch((err) => {
log.error(
"[%s] An error occurred while calling onDetached() the sender '%s': %O.",
connectionId,
sender.name,
err
);
})
);
}
// Call onDetached() on batchingReceiver so that it can gracefully close any ongoing batch operation.
const batchingReceiver = entityContext.batchingReceiver;
if (batchingReceiver && !batchingReceiver.isConnecting) {
log.error(
"[%s] calling detached on batching receiver '%s'.",
connectionId,
batchingReceiver.name
);
detachCalls.push(
batchingReceiver.onDetached(error).catch((err) => {
log.error(
"[%s] An error occurred while calling onDetached() on the batching receiver '%s': %O.",
connectionId,
batchingReceiver.name,
err
);
})
);
}
// Call onDetached() on streamingReceiver so that it can decide whether to reconnect or not
const streamingReceiver = entityContext.streamingReceiver;
if (streamingReceiver && !streamingReceiver.isConnecting) {
log.error(
"[%s] calling detached on streaming receiver '%s'.",
connectionId,
streamingReceiver.name
);
const causedByDisconnect = true;
detachCalls.push(
streamingReceiver.onDetached(error, causedByDisconnect).catch((err) => {
log.error(
"[%s] An error occurred while calling onDetached() on the streaming receiver '%s': %O.",
connectionId,
streamingReceiver.name,
err
);
})
);
}
await Promise.all(detachCalls);
};
const isManagementClientSharedWithOtherClients = (): boolean => {
for (const id of Object.keys(context.clientContexts)) {
if (
context.clientContexts[id].entityPath === entityContext.entityPath &&
context.clientContexts[id].clientId !== entityContext.clientId
) {
return true;
}
}
return false;
};
(entityContext as ClientEntityContext).close = async () => {
entityContext.isClosed = true;
if (!context.connection || !context.connection.isOpen()) {
return;
}
log.entityCtxt(
"[%s] Closing client entity context for %s: %O",
context.connectionId,
clientId
);
// Close sender
if (entityContext.sender) {
await entityContext.sender.close();
}
// Close batching receiver
if (entityContext.batchingReceiver) {
await entityContext.batchingReceiver.close();
}
// Close streaming receiver
if (entityContext.streamingReceiver) {
await entityContext.streamingReceiver.close();
}
// Close all the MessageSessions.
for (const messageSessionId of Object.keys(entityContext.messageSessions)) {
await entityContext.messageSessions[messageSessionId].close();
}
// Close the sessionManager.
if (entityContext.sessionManager) {
entityContext.sessionManager.close();
}
// Make sure that we clear the map of deferred messages
entityContext.requestResponseLockedMessages.clear();
// Delete the reference in ConnectionContext
delete context.clientContexts[clientId];
// Close the managementClient unless it is shared with other clients
if (entityContext.managementClient && !isManagementClientSharedWithOtherClients()) {
await entityContext.managementClient.close();
}
log.entityCtxt(
"[%s] Closed client entity context for %s: %O",
context.connectionId,
clientId
);
};
let managementClient = getManagementClient(context.clientContexts, entityPath);
if (!managementClient) {
const mOptions: ManagementClientOptions = {
address: options.managementClientAddress || `${entityPath}/$management`,
audience: options.managementClientAudience
};
managementClient = new ManagementClient(entityContext as ClientEntityContext, mOptions);
}
entityContext.managementClient = managementClient;
const clientEntityContext = entityContext as ClientEntityContext;
context.clientContexts[entityContext.clientId] = clientEntityContext;
log.entityCtxt("[%s] Created client entity context for %s: %O", context.connectionId, clientId);
return clientEntityContext;
}
}
// Multiple clients for the same Service Bus entity should be using the same management client.
function getManagementClient(
clients: Dictionary<ClientEntityContext>,
entityPath: string
): ManagementClient | undefined {
let result: ManagementClient | undefined;
for (const id of Object.keys(clients)) {
if (clients[id].entityPath === entityPath) {
result = clients[id].managementClient;
break;
}
}
return result;
}