From f691293768d9cb612d6996a825dad178e0ef0314 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Mon, 5 Oct 2020 18:00:05 -0700 Subject: [PATCH 01/14] Finalizing changes for adding in auto lock renewal to batching and iterators as well. --- sdk/servicebus/service-bus/CHANGELOG.md | 3 + .../service-bus/review/service-bus.api.md | 2 +- .../service-bus/src/core/autoLockRenewer.ts | 159 ++++++++----- .../service-bus/src/core/batchingReceiver.ts | 16 +- .../service-bus/src/core/messageReceiver.ts | 30 ++- .../service-bus/src/core/streamingReceiver.ts | 38 +-- sdk/servicebus/service-bus/src/models.ts | 22 +- .../service-bus/src/receivers/receiver.ts | 28 ++- .../service-bus/src/serviceBusClient.ts | 15 +- .../test/internal/abortSignal.spec.ts | 17 +- .../test/internal/autoLockRenewer.spec.ts | 109 ++++----- .../test/internal/batchingReceiver.spec.ts | 25 +- .../test/internal/receiver.spec.ts | 28 ++- .../test/internal/streamingReceiver.spec.ts | 49 +++- .../service-bus/test/renewLock.spec.ts | 225 ++++++++++-------- .../service-bus/test/retries.spec.ts | 7 +- .../test/streamingReceiver.spec.ts | 3 +- .../service-bus/test/utils/testUtils.ts | 21 +- .../service-bus/test/utils/testutils2.ts | 41 ++-- 19 files changed, 502 insertions(+), 336 deletions(-) diff --git a/sdk/servicebus/service-bus/CHANGELOG.md b/sdk/servicebus/service-bus/CHANGELOG.md index e27875e8a0b9..010e87aa4bca 100644 --- a/sdk/servicebus/service-bus/CHANGELOG.md +++ b/sdk/servicebus/service-bus/CHANGELOG.md @@ -10,6 +10,9 @@ [PR 11250](https://github.com/Azure/azure-sdk-for-js/pull/11250) - "properties" in the correlation rule filter now supports `Date`. [PR 11117](https://github.com/Azure/azure-sdk-for-js/pull/11117) +- Message locks can be auto-renewed in all receive methods (receiver.receiveMessages, receiver.subcribe + and receiver.getMessageIterator). This can be configured in options when calling `ServiceBusClient.createReceiver()`. + [PR]() ### Breaking changes diff --git a/sdk/servicebus/service-bus/review/service-bus.api.md b/sdk/servicebus/service-bus/review/service-bus.api.md index a6f81b117f98..a9c3d83631ed 100644 --- a/sdk/servicebus/service-bus/review/service-bus.api.md +++ b/sdk/servicebus/service-bus/review/service-bus.api.md @@ -126,6 +126,7 @@ export interface CreateQueueOptions extends OperationOptions { // @public export interface CreateReceiverOptions { + maxLockAutoRenewDurationInMs?: number; receiveMode?: ReceiveModeT; subQueue?: SubQueue; } @@ -190,7 +191,6 @@ export interface GetMessageIteratorOptions extends OperationOptionsBase { // @public export interface MessageHandlerOptions extends MessageHandlerOptionsBase { - maxAutoRenewLockDurationInMs?: number; } // @public diff --git a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts index 989a4c4d4ec6..9e04a39e0156 100644 --- a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts +++ b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts @@ -3,11 +3,11 @@ import { ConnectionContext } from "../connectionContext"; import { logger } from "../log"; -import { InternalReceiveMode, ServiceBusMessageImpl } from "../serviceBusMessage"; +import { ServiceBusMessageImpl } from "../serviceBusMessage"; import { logError } from "../util/errors"; import { calculateRenewAfterDuration } from "../util/utils"; import { LinkEntity } from "./linkEntity"; -import { OnError, ReceiveOptions } from "./messageReceiver"; +import { OnError } from "./messageReceiver"; /** * @internal @@ -19,27 +19,28 @@ export type RenewableMessageProperties = Readonly< // updated when we renew the lock Pick; +type MinimalLink = Pick, "name" | "logPrefix" | "entityPath">; + /** * Tracks locks for messages, renewing until a configurable duration. * * @internal * @ignore */ -export class AutoLockRenewer { +export class LockRenewer { /** - * @property {Map} _messageRenewLockTimers Maintains a map of messages for which - * the lock is automatically renewed. + * @property _messageRenewLockTimers A map of link names to individual maps for each + * link that map a message ID to its auto-renewal timer. */ - private _messageRenewLockTimers: Map = new Map< + private _messageRenewLockTimers: Map> = new Map< string, - NodeJS.Timer | undefined + Map >(); // just here for make unit testing a bit easier. private _calculateRenewAfterDuration: typeof calculateRenewAfterDuration; constructor( - private _linkEntity: Pick, "name" | "logPrefix" | "entityPath">, private _context: Pick, private _maxAutoRenewDurationInMs: number ) { @@ -56,37 +57,40 @@ export class AutoLockRenewer { * and the options.maxAutoRenewLockDurationInMs is > 0..Otherwise, returns undefined. */ static create( - linkEntity: Pick, "name" | "logPrefix" | "entityPath">, context: Pick, - options?: Pick + maxAutoRenewLockDurationInMs: number, + receiveMode: "peekLock" | "receiveAndDelete" ) { - if (options?.receiveMode === InternalReceiveMode.receiveAndDelete) { + if (receiveMode !== "peekLock") { return undefined; } - const maxAutoRenewDurationInMs = - options?.maxAutoRenewLockDurationInMs != null - ? options.maxAutoRenewLockDurationInMs - : 300 * 1000; - - if (maxAutoRenewDurationInMs <= 0) { + if (maxAutoRenewLockDurationInMs <= 0) { return undefined; } - return new AutoLockRenewer(linkEntity, context, maxAutoRenewDurationInMs); + return new LockRenewer(context, maxAutoRenewLockDurationInMs); } /** * Cancels all pending lock renewals and removes all entries from our internal cache. */ - stopAll() { + stopAll(linkEntity: MinimalLink) { logger.verbose( - `${this._linkEntity.logPrefix} Clearing message renew lock timers for all the active messages.` + `${linkEntity.logPrefix} Clearing message renew lock timers for all the active messages.` ); - for (const messageId of this._messageRenewLockTimers.keys()) { - this._stopAndRemoveById(messageId); + const messagesForLink = this._messageRenewLockTimers.get(linkEntity.name); + + if (messagesForLink == null) { + return; } + + for (const messageId of messagesForLink.keys()) { + this._stopAndRemoveById(linkEntity, messagesForLink, messageId); + } + + this._messageRenewLockTimers.delete(linkEntity.name); } /** @@ -94,9 +98,16 @@ export class AutoLockRenewer { * * @param bMessage The message whose lock renewal we will stop. */ - stop(bMessage: RenewableMessageProperties) { + stop(linkEntity: MinimalLink, bMessage: RenewableMessageProperties) { const messageId = bMessage.messageId as string; - this._stopAndRemoveById(messageId); + + const messagesForLink = this._messageRenewLockTimers.get(linkEntity.name); + + if (messagesForLink == null) { + return; + } + + this._stopAndRemoveById(linkEntity, messagesForLink, messageId); } /** @@ -104,9 +115,9 @@ export class AutoLockRenewer { * * @param bMessage The message whose lock renewal we will start. */ - start(bMessage: RenewableMessageProperties, onError: OnError) { + start(linkEntity: MinimalLink, bMessage: RenewableMessageProperties, onError: OnError) { try { - const logPrefix = this._linkEntity.logPrefix; + const logPrefix = linkEntity.logPrefix; if (bMessage.lockToken == null) { throw new Error( @@ -115,15 +126,17 @@ export class AutoLockRenewer { } const lockToken = bMessage.lockToken; + const linkMessageMap = this._getOrCreateMapForLink(linkEntity); // - We need to renew locks before they expire by looking at bMessage.lockedUntilUtc. // - This autorenewal needs to happen **NO MORE** than maxAutoRenewDurationInMs // - We should be able to clear the renewal timer when the user's message handler // is done (whether it succeeds or fails). - // Setting the messageId with undefined value in the _messageRenewockTimers Map because we + // Setting the messageId with undefined value in the linkMessageMap because we // track state by checking the presence of messageId in the map. It is removed from the map // when an attempt is made to settle the message (either by the user or by the sdk) OR // when the execution of user's message handler completes. - this._messageRenewLockTimers.set(bMessage.messageId as string, undefined); + linkMessageMap.set(bMessage.messageId as string, undefined); + logger.verbose( `${logPrefix} message with id '${ bMessage.messageId @@ -135,6 +148,7 @@ export class AutoLockRenewer { bMessage.messageId }' is: ${new Date(totalAutoLockRenewDuration).toString()}` ); + const autoRenewLockTask = (): void => { const renewalNeededToMaintainLock = // if the lock expires _after_ our max auto-renew duration there's no reason to @@ -145,7 +159,7 @@ export class AutoLockRenewer { const haventExceededMaxLockRenewalTime = Date.now() < totalAutoLockRenewDuration; if (renewalNeededToMaintainLock && haventExceededMaxLockRenewalTime) { - if (this._messageRenewLockTimers.has(bMessage.messageId as string)) { + if (linkMessageMap.has(bMessage.messageId as string)) { // TODO: We can run into problems with clock skew between the client and the server. // It would be better to calculate the duration based on the "lockDuration" property // of the queue. However, we do not have the management plane of the client ready for @@ -153,38 +167,42 @@ export class AutoLockRenewer { const amount = this._calculateRenewAfterDuration(bMessage.lockedUntilUtc!); logger.verbose( - `${logPrefix} Sleeping for %d milliseconds while renewing the lock for message with id '${bMessage.messageId}' is: ${amount}` + `${logPrefix} Sleeping for ${amount} milliseconds while renewing the lock for message with id '${bMessage.messageId}'` ); // Setting the value of the messageId to the actual timer. This will be cleared when // an attempt is made to settle the message (either by the user or by the sdk) OR // when the execution of user's message handler completes. - this._messageRenewLockTimers.set( - bMessage.messageId as string, - setTimeout(async () => { - try { - logger.verbose( - `${logPrefix} Attempting to renew the lock for message with id '${bMessage.messageId}'.` - ); - - bMessage.lockedUntilUtc = await this._context - .getManagementClient(this._linkEntity.entityPath) - .renewLock(lockToken, { - associatedLinkName: this._linkEntity.name - }); - logger.verbose( - `${logPrefix} Successfully renewed the lock for message with id '${bMessage.messageId}'. Starting next auto-lock-renew cycle for message.` - ); - - autoRenewLockTask(); - } catch (err) { - logError( - err, - `${logPrefix} An error occurred while auto renewing the message lock '${bMessage.lockToken}' for message with id '${bMessage.messageId}'` - ); - onError(err); - } - }, amount) - ); + const autoRenewTimer = setTimeout(async () => { + try { + logger.verbose( + `${logPrefix} Attempting to renew the lock for message with id '${bMessage.messageId}'.` + ); + + bMessage.lockedUntilUtc = await this._context + .getManagementClient(linkEntity.entityPath) + .renewLock(lockToken, { + associatedLinkName: linkEntity.name + }); + logger.verbose( + `${logPrefix} Successfully renewed the lock for message with id '${bMessage.messageId}'. Starting next auto-lock-renew cycle for message.` + ); + + autoRenewLockTask(); + } catch (err) { + logError( + err, + `${logPrefix} An error occurred while auto renewing the message lock '${bMessage.lockToken}' for message with id '${bMessage.messageId}'` + ); + onError(err); + } + }, amount); + + // Prevent the active Timer from keeping the Node.js event loop active. + if (typeof autoRenewTimer.unref === "function") { + autoRenewTimer.unref(); + } + + linkMessageMap.set(bMessage.messageId as string, autoRenewTimer); } else { logger.verbose( `${logPrefix} Looks like the message lock renew timer has already been cleared for message with id '${bMessage.messageId}'.` @@ -199,7 +217,7 @@ export class AutoLockRenewer { }'. Hence we will stop the autoLockRenewTask.` ); - this.stop(bMessage); + this.stop(linkEntity, bMessage); } }; @@ -210,19 +228,34 @@ export class AutoLockRenewer { } } - private _stopAndRemoveById(messageId: string | undefined): void { + private _getOrCreateMapForLink(linkEntity: MinimalLink): Map { + if (!this._messageRenewLockTimers.has(linkEntity.name)) { + this._messageRenewLockTimers.set( + linkEntity.name, + new Map() + ); + } + + return this._messageRenewLockTimers.get(linkEntity.name)!; + } + + private _stopAndRemoveById( + linkEntity: MinimalLink, + linkMessageMap: Map, + messageId: string | undefined + ): void { if (messageId == null) { throw new Error("Failed to stop auto lock renewal - no message ID"); } // TODO: messageId doesn't actually need to be unique. Perhaps we should use lockToken // instead? - if (this._messageRenewLockTimers.has(messageId)) { - clearTimeout(this._messageRenewLockTimers.get(messageId) as NodeJS.Timer); + if (linkMessageMap.has(messageId)) { + clearTimeout(linkMessageMap.get(messageId) as NodeJS.Timer); logger.verbose( - `${this._linkEntity.logPrefix} Cleared the message renew lock timer for message with id '${messageId}'.` + `${linkEntity.logPrefix} Cleared the message renew lock timer for message with id '${messageId}'.` ); - this._messageRenewLockTimers.delete(messageId); + linkMessageMap.delete(messageId); } } } diff --git a/sdk/servicebus/service-bus/src/core/batchingReceiver.ts b/sdk/servicebus/service-bus/src/core/batchingReceiver.ts index fcd52ee842af..aae238925309 100644 --- a/sdk/servicebus/service-bus/src/core/batchingReceiver.ts +++ b/sdk/servicebus/service-bus/src/core/batchingReceiver.ts @@ -35,7 +35,7 @@ export class BatchingReceiver extends MessageReceiver { * @param {ClientEntityContext} context The client entity context. * @param {ReceiveOptions} [options] Options for how you'd like to connect. */ - constructor(context: ConnectionContext, entityPath: string, options?: ReceiveOptions) { + constructor(context: ConnectionContext, entityPath: string, options: ReceiveOptions) { super(context, entityPath, "br", options); this._batchingReceiverLite = new BatchingReceiverLite( @@ -118,12 +118,22 @@ export class BatchingReceiver extends MessageReceiver { this.name ); - return await this._batchingReceiverLite.receiveMessages({ + const messages = await this._batchingReceiverLite.receiveMessages({ maxMessageCount, maxWaitTimeInMs, maxTimeAfterFirstMessageInMs, userAbortSignal }); + + if (this._lockRenewer) { + for (const message of messages) { + this._lockRenewer.start(this, message, (error) => { + logError(error, `${this.logPrefix} Failed to renew lock for message.`); + }); + } + } + + return messages; } catch (error) { logError( error, @@ -139,7 +149,7 @@ export class BatchingReceiver extends MessageReceiver { static create( context: ConnectionContext, entityPath: string, - options?: ReceiveOptions + options: ReceiveOptions ): BatchingReceiver { throwErrorIfConnectionClosed(context); const bReceiver = new BatchingReceiver(context, entityPath, options); diff --git a/sdk/servicebus/service-bus/src/core/messageReceiver.ts b/sdk/servicebus/service-bus/src/core/messageReceiver.ts index 3bd2b9904021..7c10f642c34f 100644 --- a/sdk/servicebus/service-bus/src/core/messageReceiver.ts +++ b/sdk/servicebus/service-bus/src/core/messageReceiver.ts @@ -19,7 +19,7 @@ import { DispositionStatusOptions } from "./managementClient"; import { AbortSignalLike } from "@azure/core-http"; import { onMessageSettled, DeferredPromiseAndTimer } from "./shared"; import { logError } from "../util/errors"; -import { AutoLockRenewer } from "./autoLockRenewer"; +import { LockRenewer } from "./autoLockRenewer"; /** * @internal @@ -45,13 +45,19 @@ export interface OnAmqpEventAsPromise extends OnAmqpEvent { export interface ReceiveOptions extends MessageHandlerOptions { /** * @property {number} [receiveMode] The mode in which messages should be received. - * Default: ReceiveMode.peekLock */ - receiveMode?: InternalReceiveMode; + receiveMode: InternalReceiveMode; /** * Retry policy options that determine the mode, number of retries, retry interval etc. */ retryOptions?: RetryOptions; + + /** + * A LockAutoRenewer that will automatically renew locks based on user specified interval. + * This will be set if the user has chosen peekLock mode _and_ they've set a positive + * maxAutoRenewLockDurationInMs value when they created their receiver. + */ + lockRenewer: LockRenewer | undefined; } /** @@ -120,31 +126,31 @@ export abstract class MessageReceiver extends LinkEntity { * inside _onAmqpError. */ protected _onError?: OnError; + /** - * An AutoLockRenewer. This is undefined unless the user has activated autolock renewal - * via ReceiveOptions. + * A lock renewer that handles message lock auto-renewal. This is undefined unless the user + * has activated autolock renewal via ReceiveOptions. A single auto lock renewer is shared + * for all links for a `ServiceBusReceiver` instance. */ - protected _autolockRenewer: AutoLockRenewer | undefined; + protected _lockRenewer: LockRenewer | undefined; constructor( context: ConnectionContext, entityPath: string, receiverType: ReceiverType, - options?: Omit + options: Omit ) { super(entityPath, entityPath, context, receiverType, { address: entityPath, audience: `${context.config.endpoint}${entityPath}` }); - if (!options) options = {}; this.receiverType = receiverType; this.receiveMode = options.receiveMode || InternalReceiveMode.peekLock; // If explicitly set to false then autoComplete is false else true (default). this.autoComplete = options.autoComplete === false ? options.autoComplete : true; - - this._autolockRenewer = AutoLockRenewer.create(this, this._context, options); + this._lockRenewer = options.lockRenewer; } /** @@ -231,7 +237,7 @@ export abstract class MessageReceiver extends LinkEntity { * @return {Promise} Promise. */ async close(): Promise { - this._autolockRenewer?.stopAll(); + this._lockRenewer?.stopAll(this); await super.close(); } @@ -251,7 +257,7 @@ export abstract class MessageReceiver extends LinkEntity { if (operation.match(/^(complete|abandon|defer|deadletter)$/) == null) { return reject(new Error(`operation: '${operation}' is not a valid operation.`)); } - this._autolockRenewer?.stop(message); + this._lockRenewer?.stop(this, message); const delivery = message.delivery; const timer = setTimeout(() => { this._deliveryDispositionMap.delete(delivery.id); diff --git a/sdk/servicebus/service-bus/src/core/streamingReceiver.ts b/sdk/servicebus/service-bus/src/core/streamingReceiver.ts index 9146a8fc4632..5f0d35c214f9 100644 --- a/sdk/servicebus/service-bus/src/core/streamingReceiver.ts +++ b/sdk/servicebus/service-bus/src/core/streamingReceiver.ts @@ -29,6 +29,23 @@ import { AmqpError, EventContext, isAmqpError, OnAmqpEvent } from "rhea-promise" import { InternalReceiveMode, ServiceBusMessageImpl } from "../serviceBusMessage"; import { AbortSignalLike } from "@azure/abort-controller"; +/** + * @internal + * @ignore + */ +export interface CreateStreamingReceiverOptions + extends ReceiveOptions, + Pick { + /** + * Used for mocking/stubbing in tests. + */ + _createStreamingReceiverStubForTests?: ( + context: ConnectionContext, + options?: ReceiveOptions + ) => StreamingReceiver; + cachedStreamingReceiver?: StreamingReceiver; +} + /** * @internal * @ignore @@ -106,7 +123,7 @@ export class StreamingReceiver extends MessageReceiver { * @param {ClientEntityContext} context The client entity context. * @param {ReceiveOptions} [options] Options for how you'd like to connect. */ - constructor(context: ConnectionContext, entityPath: string, options?: ReceiveOptions) { + constructor(context: ConnectionContext, entityPath: string, options: ReceiveOptions) { super(context, entityPath, "sr", options); if (typeof options?.maxConcurrentCalls === "number" && options?.maxConcurrentCalls > 0) { @@ -127,7 +144,7 @@ export class StreamingReceiver extends MessageReceiver { receiverError ); - this._autolockRenewer?.stopAll(); + this._lockRenewer?.stopAll(this); if (receiver && !receiver.isItselfClosed()) { await this.onDetached(receiverError); @@ -154,7 +171,7 @@ export class StreamingReceiver extends MessageReceiver { sessionError ); - this._autolockRenewer?.stopAll(); + this._lockRenewer?.stopAll(this); if (receiver && !receiver.isSessionItselfClosed()) { await this.onDetached(sessionError); @@ -258,7 +275,7 @@ export class StreamingReceiver extends MessageReceiver { this.receiveMode ); - this._autolockRenewer?.start(bMessage, (err) => { + this._lockRenewer?.start(this, bMessage, (err) => { if (this._onError) { this._onError(err); } @@ -266,7 +283,6 @@ export class StreamingReceiver extends MessageReceiver { try { await this._onMessage(bMessage); - this._autolockRenewer?.stop(bMessage); } catch (err) { // This ensures we call users' error handler when users' message handler throws. if (!isAmqpError(err)) { @@ -284,7 +300,7 @@ export class StreamingReceiver extends MessageReceiver { // Do not want renewLock to happen unnecessarily, while abandoning the message. Hence, // doing this here. Otherwise, this should be done in finally. - this._autolockRenewer?.stop(bMessage); + this._lockRenewer?.stop(this, bMessage); const error = translate(err) as MessagingError; // Nothing much to do if user's message handler throws. Let us try abandoning the message. if ( @@ -548,17 +564,9 @@ export class StreamingReceiver extends MessageReceiver { static async create( context: ConnectionContext, entityPath: string, - options?: ReceiveOptions & - Pick & { - _createStreamingReceiverStubForTests?: ( - context: ConnectionContext, - options?: ReceiveOptions - ) => StreamingReceiver; - cachedStreamingReceiver?: StreamingReceiver; - } + options: CreateStreamingReceiverOptions ): Promise { throwErrorIfConnectionClosed(context); - if (!options) options = {}; if (options.autoComplete == null) options.autoComplete = true; let sReceiver: StreamingReceiver; diff --git a/sdk/servicebus/service-bus/src/models.ts b/sdk/servicebus/service-bus/src/models.ts index fc88eafb442f..06cb8cc34a57 100644 --- a/sdk/servicebus/service-bus/src/models.ts +++ b/sdk/servicebus/service-bus/src/models.ts @@ -80,6 +80,16 @@ export interface CreateReceiverOptions { * see https://docs.microsoft.com/azure/service-bus-messaging/service-bus-dead-letter-queues */ subQueue?: SubQueue; + + /** + * The maximum duration in milliseconds until which the lock on the message will be renewed + * by the sdk automatically. This auto renewal stops once the message is settled or once the user + * provided onMessage handler completes ite execution. + * + * - **Default**: `300 * 1000` milliseconds (5 minutes). + * - **To disable autolock renewal**, set this to `0`. + */ + maxLockAutoRenewDurationInMs?: number; } /** @@ -152,17 +162,7 @@ export interface MessageHandlerOptionsBase extends OperationOptionsBase { * Describes the options passed to `registerMessageHandler` method when receiving messages from a * Queue/Subscription which does not have sessions enabled. */ -export interface MessageHandlerOptions extends MessageHandlerOptionsBase { - /** - * @property The maximum duration in milliseconds until which the lock on the message will be renewed - * by the sdk automatically. This auto renewal stops once the message is settled or once the user - * provided onMessage handler completes ite execution. - * - * - **Default**: `300 * 1000` milliseconds (5 minutes). - * - **To disable autolock renewal**, set this to `0`. - */ - maxAutoRenewLockDurationInMs?: number; -} +export interface MessageHandlerOptions extends MessageHandlerOptionsBase {} /** * Describes the options passed to the `acceptSession` and `acceptNextSession` methods diff --git a/sdk/servicebus/service-bus/src/receivers/receiver.ts b/sdk/servicebus/service-bus/src/receivers/receiver.ts index 85f9d23057b8..7b6ae2d13880 100644 --- a/sdk/servicebus/service-bus/src/receivers/receiver.ts +++ b/sdk/servicebus/service-bus/src/receivers/receiver.ts @@ -21,7 +21,7 @@ import { throwTypeErrorIfParameterNotLong } from "../util/errors"; import { OnError, OnMessage, ReceiveOptions } from "../core/messageReceiver"; -import { StreamingReceiver } from "../core/streamingReceiver"; +import { CreateStreamingReceiverOptions, StreamingReceiver } from "../core/streamingReceiver"; import { BatchingReceiver } from "../core/batchingReceiver"; import { assertValidMessageHandlers, getMessageIterator, wrapProcessErrorHandler } from "./shared"; import { convertToInternalReceiveMode } from "../constructorHelpers"; @@ -29,6 +29,7 @@ import Long from "long"; import { ServiceBusReceivedMessageWithLock, ServiceBusMessageImpl } from "../serviceBusMessage"; import { Constants, RetryConfig, RetryOperationType, RetryOptions, retry } from "@azure/core-amqp"; import "@azure/core-asynciterator-polyfill"; +import { LockRenewer } from "../core/autoLockRenewer"; /** * A receiver that does not handle sessions. @@ -155,6 +156,7 @@ export class ServiceBusReceiverImpl< * Instance of the StreamingReceiver class to use to receive messages in push mode. */ private _streamingReceiver?: StreamingReceiver; + private _lockRenewer: LockRenewer | undefined; /** * @throws Error if the underlying connection is closed. @@ -163,10 +165,16 @@ export class ServiceBusReceiverImpl< private _context: ConnectionContext, public entityPath: string, public receiveMode: "peekLock" | "receiveAndDelete", + maxAutoRenewLockDurationInMs: number, retryOptions: RetryOptions = {} ) { throwErrorIfConnectionClosed(_context); this._retryOptions = retryOptions; + this._lockRenewer = LockRenewer.create( + this._context, + maxAutoRenewLockDurationInMs, + receiveMode + ); } private _throwIfAlreadyReceiving(): void { @@ -235,7 +243,8 @@ export class ServiceBusReceiverImpl< ...options, receiveMode: convertToInternalReceiveMode(this.receiveMode), retryOptions: this._retryOptions, - cachedStreamingReceiver: this._streamingReceiver + cachedStreamingReceiver: this._streamingReceiver, + lockRenewer: this._lockRenewer }) .then(async (sReceiver) => { if (!sReceiver) { @@ -264,15 +273,7 @@ export class ServiceBusReceiverImpl< private _createStreamingReceiver( context: ConnectionContext, entityPath: string, - options?: ReceiveOptions & - Pick & { - createStreamingReceiver?: ( - context: ConnectionContext, - entityPath: string, - options?: ReceiveOptions - ) => StreamingReceiver; - cachedStreamingReceiver?: StreamingReceiver; - } + options: CreateStreamingReceiverOptions ): Promise { return StreamingReceiver.create(context, entityPath, options); } @@ -292,7 +293,8 @@ export class ServiceBusReceiverImpl< if (!this._batchingReceiver || !this._context.messageReceivers[this._batchingReceiver.name]) { const options: ReceiveOptions = { maxConcurrentCalls: 0, - receiveMode: convertToInternalReceiveMode(this.receiveMode) + receiveMode: convertToInternalReceiveMode(this.receiveMode), + lockRenewer: this._lockRenewer }; this._batchingReceiver = this._createBatchingReceiver( this._context, @@ -498,7 +500,7 @@ export class ServiceBusReceiverImpl< private _createBatchingReceiver( context: ConnectionContext, entityPath: string, - options?: ReceiveOptions + options: ReceiveOptions ): BatchingReceiver { return BatchingReceiver.create(context, entityPath, options); } diff --git a/sdk/servicebus/service-bus/src/serviceBusClient.ts b/sdk/servicebus/service-bus/src/serviceBusClient.ts index c06cb2408a3d..e6103344331d 100644 --- a/sdk/servicebus/service-bus/src/serviceBusClient.ts +++ b/sdk/servicebus/service-bus/src/serviceBusClient.ts @@ -236,11 +236,17 @@ export class ServiceBusClient { } } + const maxLockAutoRenewDurationInMs = + options?.maxLockAutoRenewDurationInMs != null + ? options.maxLockAutoRenewDurationInMs + : 5 * 60 * 1000; + if (receiveMode === "peekLock") { return new ServiceBusReceiverImpl( this._connectionContext, entityPathWithSubQueue, receiveMode, + maxLockAutoRenewDurationInMs, this._clientOptions.retryOptions ); } else { @@ -248,6 +254,7 @@ export class ServiceBusClient { this._connectionContext, entityPathWithSubQueue, receiveMode, + maxLockAutoRenewDurationInMs, this._clientOptions.retryOptions ); } @@ -353,9 +360,7 @@ export class ServiceBusClient { | AcceptSessionOptions<"peekLock"> | AcceptSessionOptions<"receiveAndDelete"> | string, - options4?: - | AcceptSessionOptions<"peekLock"> - | AcceptSessionOptions<"receiveAndDelete"> + options4?: AcceptSessionOptions<"peekLock"> | AcceptSessionOptions<"receiveAndDelete"> ): Promise< | ServiceBusSessionReceiver | ServiceBusSessionReceiver @@ -514,9 +519,7 @@ export class ServiceBusClient { | AcceptSessionOptions<"peekLock"> | AcceptSessionOptions<"receiveAndDelete"> | string, - options3?: - | AcceptSessionOptions<"peekLock"> - | AcceptSessionOptions<"receiveAndDelete"> + options3?: AcceptSessionOptions<"peekLock"> | AcceptSessionOptions<"receiveAndDelete"> ): Promise< | ServiceBusSessionReceiver | ServiceBusSessionReceiver diff --git a/sdk/servicebus/service-bus/test/internal/abortSignal.spec.ts b/sdk/servicebus/service-bus/test/internal/abortSignal.spec.ts index 2c426e438569..e7dc0702a3c5 100644 --- a/sdk/servicebus/service-bus/test/internal/abortSignal.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/abortSignal.spec.ts @@ -24,8 +24,14 @@ import { isLinkLocked } from "../utils/misc"; import { ServiceBusSessionReceiverImpl } from "../../src/receivers/sessionReceiver"; import { ServiceBusReceiverImpl } from "../../src/receivers/receiver"; import { MessageSession } from "../../src/session/messageSession"; +import { InternalReceiveMode } from "../../src/serviceBusMessage"; describe("AbortSignal", () => { + const defaultOptions = { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + }; + const testMessageThatDoesntMatter = { body: "doesn't matter" }; @@ -253,7 +259,8 @@ describe("AbortSignal", () => { it("...before first async call", async () => { const messageReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(messageReceiver); @@ -273,7 +280,8 @@ describe("AbortSignal", () => { it("...after negotiateClaim", async () => { const messageReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(messageReceiver); @@ -304,7 +312,7 @@ describe("AbortSignal", () => { isAborted = true; } }); - const messageReceiver = new StreamingReceiver(fakeContext, "fakeEntityPath"); + const messageReceiver = new StreamingReceiver(fakeContext, "fakeEntityPath", defaultOptions); closeables.push(messageReceiver); messageReceiver["_negotiateClaim"] = async () => {}; @@ -366,7 +374,8 @@ describe("AbortSignal", () => { const receiver = new ServiceBusReceiverImpl( createConnectionContextForTests(), "entityPath", - "peekLock" + "peekLock", + 1 ); try { diff --git a/sdk/servicebus/service-bus/test/internal/autoLockRenewer.spec.ts b/sdk/servicebus/service-bus/test/internal/autoLockRenewer.spec.ts index f63b0bcf3d46..b9413160264a 100644 --- a/sdk/servicebus/service-bus/test/internal/autoLockRenewer.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/autoLockRenewer.spec.ts @@ -7,15 +7,14 @@ import chaiAsPromised from "chai-as-promised"; chai.use(chaiAsPromised); const assert = chai.assert; import * as sinon from "sinon"; -import { AutoLockRenewer } from "../../src/core/autoLockRenewer"; +import { LockRenewer } from "../../src/core/autoLockRenewer"; import { ManagementClient, SendManagementRequestOptions } from "../../src/core/managementClient"; -import { InternalReceiveMode } from "../../src/serviceBusMessage"; import { getPromiseResolverForTest } from "./unittestUtils"; describe("autoLockRenewer unit tests", () => { let clock: ReturnType; - let autoLockRenewer: AutoLockRenewer; + let autoLockRenewer: LockRenewer; let renewLockSpy: sinon.SinonSpy< Parameters, @@ -30,6 +29,12 @@ describe("autoLockRenewer unit tests", () => { msToNextRenewal: 5 }; + const testLinkEntity = { + name: "linkName", + logPrefix: "this is my log prefix", + entityPath: "entity path" + }; + let stopTimerPromise: Promise; beforeEach(() => { @@ -48,22 +53,15 @@ describe("autoLockRenewer unit tests", () => { renewLockSpy = sinon.spy(managementClient, "renewLock"); onErrorFake = sinon.fake(async (_err: Error | MessagingError) => {}); - autoLockRenewer = AutoLockRenewer.create( - { - name: "linkName", - logPrefix: "this is my log prefix", - entityPath: "entity path" - }, + autoLockRenewer = LockRenewer.create( { getManagementClient: (entityPath) => { assert.equal(entityPath, "entity path"); return managementClient; } }, - { - maxAutoRenewLockDurationInMs: limits.maxAdditionalTimeToRenewLock, - receiveMode: InternalReceiveMode.peekLock - } + limits.maxAdditionalTimeToRenewLock, + "peekLock" )!; // always just start the next auto-renew timer after 5 milliseconds to keep things simple. @@ -74,13 +72,29 @@ describe("autoLockRenewer unit tests", () => { const origStop = autoLockRenewer["stop"].bind(autoLockRenewer); - autoLockRenewer["stop"] = (message) => { - origStop(message); + autoLockRenewer["stop"] = (linkEntity, message) => { + origStop(linkEntity, message); stopTimerResolve(); }; }); afterEach(() => { + // each test should properly end "clean" as far as removing any + // message timers. + let lockRenewalTimersTotal = 0; + for (const value of autoLockRenewer["_messageRenewLockTimers"].values()) { + lockRenewalTimersTotal += value.size; + } + + assert.equal( + lockRenewalTimersTotal, + 0, + "Should be no active lock timers after test have completed." + ); + + // the per-link map is not cleaned up automatically - you must + // stopAll() to remove it. + autoLockRenewer.stopAll(testLinkEntity); assert.equal( autoLockRenewer["_messageRenewLockTimers"].size, 0, @@ -92,6 +106,7 @@ describe("autoLockRenewer unit tests", () => { it("standard renewal", async () => { autoLockRenewer.start( + testLinkEntity, { lockToken: "lock token", lockedUntilUtc: new Date(), @@ -103,7 +118,7 @@ describe("autoLockRenewer unit tests", () => { clock.tick(limits.msToNextRenewal - 1); // right before the renew timer would run assert.exists( - autoLockRenewer["_messageRenewLockTimers"].get("message id"), + autoLockRenewer["_messageRenewLockTimers"].get(testLinkEntity.name)?.get("message id"), "auto-renew timer should be set up" ); @@ -123,8 +138,21 @@ describe("autoLockRenewer unit tests", () => { assert.isFalse(onErrorFake.called, "no errors"); }); + it("delete multiple times", () => { + // no lock renewal for this message + autoLockRenewer.stop(testLinkEntity, { + messageId: "hello" + }); + + // no locks have been setup + autoLockRenewer.stopAll(testLinkEntity); + + assert.isEmpty(autoLockRenewer["_messageRenewLockTimers"]); + }); + it("renewal timer not scheduled: message is already locked for longer than our renewal would extend it", () => { autoLockRenewer.start( + testLinkEntity, { lockToken: "lock token", // this date exceeds the max time we would renew for so we don't need to do anything. @@ -144,6 +172,7 @@ describe("autoLockRenewer unit tests", () => { it("renewal timer is not (re)scheduled: the current date has passed our max lock renewal time", async () => { autoLockRenewer.start( + testLinkEntity, { lockToken: "lock token", lockedUntilUtc: new Date(), @@ -175,6 +204,7 @@ describe("autoLockRenewer unit tests", () => { it("invalid message can't renew", () => { autoLockRenewer.start( + testLinkEntity, { messageId: "my message id" }, @@ -198,17 +228,10 @@ describe("autoLockRenewer unit tests", () => { }; it("doesn't support receiveAndDelete mode", () => { - const autoLockRenewer = AutoLockRenewer.create( - { - name: "linkName", - logPrefix: "this is my log prefix", - entityPath: "entity path" - }, + const autoLockRenewer = LockRenewer.create( unusedMgmtClient, - { - maxAutoRenewLockDurationInMs: 1, // this is okay - receiveMode: InternalReceiveMode.receiveAndDelete // this is not okay - there aren't any locks to renew in receiveAndDelete mode. - } + 1, // this is okay, + "receiveAndDelete" // this is not okay - there aren't any locks to renew in receiveAndDelete mode. ); assert.notExists( @@ -219,17 +242,10 @@ describe("autoLockRenewer unit tests", () => { [0, -1].forEach((invalidMaxAutoRenewLockDurationInMs) => { it(`Invalid maxAutoRenewLockDurationInMs duration: ${invalidMaxAutoRenewLockDurationInMs}`, () => { - const autoLockRenewer = AutoLockRenewer.create( - { - name: "linkName", - logPrefix: "this is my log prefix", - entityPath: "entity path" - }, + const autoLockRenewer = LockRenewer.create( unusedMgmtClient, - { - receiveMode: InternalReceiveMode.peekLock, // this is okay - maxAutoRenewLockDurationInMs: invalidMaxAutoRenewLockDurationInMs - } + invalidMaxAutoRenewLockDurationInMs, + "peekLock" // this is okay ); assert.notExists( @@ -238,26 +254,5 @@ describe("autoLockRenewer unit tests", () => { ); }); }); - - it(`maxAutoRenewLockDurationInMs of undefined becomes the default (5 minutes)`, () => { - const autoLockRenewer = AutoLockRenewer.create( - { - name: "linkName", - logPrefix: "this is my log prefix", - entityPath: "entity path" - }, - unusedMgmtClient, - { - receiveMode: InternalReceiveMode.peekLock - } - ); - - assert.exists(autoLockRenewer); - assert.equal( - autoLockRenewer!["_maxAutoRenewDurationInMs"], - 1000 * 5 * 60, - "By default our max auto renew lock duration is 5 minutes" - ); - }); }); }); diff --git a/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts b/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts index acb8360a8dde..df8c24aa2701 100644 --- a/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts @@ -30,6 +30,7 @@ import { StandardAbortMessage } from "../../src/util/utils"; import { OnAmqpEventAsPromise } from "../../src/core/messageReceiver"; import { ConnectionContext } from "../../src/connectionContext"; import { ServiceBusReceiverImpl } from "../../src/receivers/receiver"; +import { LockRenewer } from "../../src/core/autoLockRenewer"; describe("BatchingReceiver unit tests", () => { let closeables: { close(): Promise }[]; @@ -52,7 +53,8 @@ describe("BatchingReceiver unit tests", () => { const receiver = new ServiceBusReceiverImpl( createConnectionContextForTests(), "fakeEntityPath", - "peekLock" + "peekLock", + 1 ); let wasCalled = false; @@ -84,7 +86,8 @@ describe("BatchingReceiver unit tests", () => { abortController.abort(); const receiver = new BatchingReceiver(createConnectionContextForTests(), "fakeEntityPath", { - receiveMode: InternalReceiveMode.peekLock + receiveMode: InternalReceiveMode.peekLock, + lockRenewer: undefined }); try { @@ -100,7 +103,8 @@ describe("BatchingReceiver unit tests", () => { const abortController = new AbortController(); const receiver = new BatchingReceiver(createConnectionContextForTests(), "fakeEntityPath", { - receiveMode: InternalReceiveMode.peekLock + receiveMode: InternalReceiveMode.peekLock, + lockRenewer: undefined }); closeables.push(receiver); @@ -185,7 +189,8 @@ describe("BatchingReceiver unit tests", () => { createConnectionContextForTests(), "dummyEntityPath", { - receiveMode: lockMode + receiveMode: lockMode, + lockRenewer: undefined } ); closeables.push(receiver); @@ -218,7 +223,8 @@ describe("BatchingReceiver unit tests", () => { createConnectionContextForTests(), "dummyEntityPath", { - receiveMode: lockMode + receiveMode: lockMode, + lockRenewer: undefined } ); closeables.push(receiver); @@ -249,7 +255,8 @@ describe("BatchingReceiver unit tests", () => { createConnectionContextForTests(), "dummyEntityPath", { - receiveMode: lockMode + receiveMode: lockMode, + lockRenewer: undefined } ); closeables.push(receiver); @@ -300,7 +307,8 @@ describe("BatchingReceiver unit tests", () => { createConnectionContextForTests(), "dummyEntityPath", { - receiveMode: lockMode + receiveMode: lockMode, + lockRenewer: undefined } ); closeables.push(receiver); @@ -356,7 +364,8 @@ describe("BatchingReceiver unit tests", () => { createConnectionContextForTests(), "dummyEntityPath", { - receiveMode: lockMode + receiveMode: lockMode, + lockRenewer: undefined } ); closeables.push(receiver); diff --git a/sdk/servicebus/service-bus/test/internal/receiver.spec.ts b/sdk/servicebus/service-bus/test/internal/receiver.spec.ts index 2b71ca65fb0b..b23d134f929e 100644 --- a/sdk/servicebus/service-bus/test/internal/receiver.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/receiver.spec.ts @@ -20,13 +20,18 @@ import { createAbortSignalForTest } from "../utils/abortSignalTestUtils"; import { AbortSignalLike } from "@azure/abort-controller"; import { ServiceBusSessionReceiverImpl } from "../../src/receivers/sessionReceiver"; import { MessageSession } from "../../src/session/messageSession"; +import { InternalReceiveMode } from "../../src/serviceBusMessage"; describe("Receiver unit tests", () => { describe("init() and close() interactions", () => { it("close() called just after init() but before the next step", async () => { const batchingReceiver = new BatchingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + } ); let initWasCalled = false; @@ -47,7 +52,11 @@ describe("Receiver unit tests", () => { it("message receiver init() bails out early if object is closed()", async () => { const messageReceiver2 = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + } ); await messageReceiver2.close(); @@ -91,7 +100,8 @@ describe("Receiver unit tests", () => { } }), "fakeEntityPath", - "peekLock" + "peekLock", + 0 ); const subscription = await subscribeAndWaitForInitialize(receiverImpl); @@ -119,7 +129,8 @@ describe("Receiver unit tests", () => { const receiverImpl = new ServiceBusReceiverImpl( createConnectionContextForTests(), "fakeEntityPath", - "peekLock" + "peekLock", + 1 ); const subscription = await subscribeAndWaitForInitialize(receiverImpl); @@ -156,7 +167,8 @@ describe("Receiver unit tests", () => { } }), "fakeEntityPath", - "peekLock" + "peekLock", + 1 ); const subscription = await subscribeAndWaitForInitialize(receiverImpl); @@ -186,7 +198,8 @@ describe("Receiver unit tests", () => { const receiverImpl = new ServiceBusReceiverImpl( createConnectionContextForTests(), "fakeEntityPath", - "peekLock" + "peekLock", + 1 ); const abortSignal = { @@ -250,7 +263,8 @@ describe("Receiver unit tests", () => { const impl = new ServiceBusReceiverImpl( createConnectionContextForTests(), "entity path", - "peekLock" + "peekLock", + 1 ); const abortSignal = createAbortSignalForTest(true); diff --git a/sdk/servicebus/service-bus/test/internal/streamingReceiver.spec.ts b/sdk/servicebus/service-bus/test/internal/streamingReceiver.spec.ts index c8532c8812a6..4991d7445006 100644 --- a/sdk/servicebus/service-bus/test/internal/streamingReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/streamingReceiver.spec.ts @@ -11,10 +11,16 @@ import { OperationOptions } from "../../src"; import { StreamingReceiver } from "../../src/core/streamingReceiver"; import { AbortController, AbortSignalLike } from "@azure/abort-controller"; import sinon from "sinon"; +import { InternalReceiveMode } from "../../src/serviceBusMessage"; chai.use(chaiAsPromised); const assert = chai.assert; describe("StreamingReceiver unit tests", () => { + const defaultOptions = { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + }; + let closeables: { close(): Promise }[]; beforeEach(() => { @@ -31,7 +37,8 @@ describe("StreamingReceiver unit tests", () => { it("off by default", async () => { const streamingReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(streamingReceiver); assert.isFalse(streamingReceiver.isReceivingMessages); @@ -47,7 +54,8 @@ describe("StreamingReceiver unit tests", () => { createConnectionContextForTests(), "fakeEntityPath", { - maxConcurrentCalls: 101 + maxConcurrentCalls: 101, + ...defaultOptions } ); closeables.push(streamingReceiver); @@ -107,7 +115,8 @@ describe("StreamingReceiver unit tests", () => { it("isReceivingMessages set to false by close()'ing", async () => { const streamingReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(streamingReceiver); @@ -127,7 +136,8 @@ describe("StreamingReceiver unit tests", () => { it("isReceivingMessages set to false by calling stopReceivingMessages()", async () => { const streamingReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(streamingReceiver); @@ -147,7 +157,8 @@ describe("StreamingReceiver unit tests", () => { it("isReceivingMessages set to false by calling onDetach and init fails", async () => { const streamingReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(streamingReceiver); @@ -171,7 +182,8 @@ describe("StreamingReceiver unit tests", () => { it("isReceivingMessages is set to true if onDetach succeeds in reconnecting", async () => { const streamingReceiver = new StreamingReceiver( createConnectionContextForTests(), - "fakeEntityPath" + "fakeEntityPath", + defaultOptions ); closeables.push(streamingReceiver); @@ -198,7 +210,10 @@ describe("StreamingReceiver unit tests", () => { it("create() with an existing receiver and that receiver is open()", async () => { const context = createConnectionContextForTests(); - const existingStreamingReceiver = new StreamingReceiver(context, "fakeEntityPath"); + const existingStreamingReceiver = new StreamingReceiver(context, "fakeEntityPath", { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + }); closeables.push(existingStreamingReceiver); await existingStreamingReceiver.init(false); @@ -208,7 +223,9 @@ describe("StreamingReceiver unit tests", () => { const spy = sinon.spy(existingStreamingReceiver, "init"); const newStreamingReceiver = await StreamingReceiver.create(context, "fakeEntityPath", { - cachedStreamingReceiver: existingStreamingReceiver + cachedStreamingReceiver: existingStreamingReceiver, + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock }); assert.isTrue(spy.called, "We do still call init() on the receiver"); @@ -229,7 +246,10 @@ describe("StreamingReceiver unit tests", () => { it("create() with an existing receiver and that receiver is NOT open()", async () => { const context = createConnectionContextForTests(); - const existingStreamingReceiver = new StreamingReceiver(context, "fakeEntityPath"); + const existingStreamingReceiver = new StreamingReceiver(context, "fakeEntityPath", { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + }); closeables.push(existingStreamingReceiver); await existingStreamingReceiver.init(false); @@ -247,7 +267,9 @@ describe("StreamingReceiver unit tests", () => { ); const newStreamingReceiver = await StreamingReceiver.create(context, "fakeEntityPath", { - cachedStreamingReceiver: existingStreamingReceiver + cachedStreamingReceiver: existingStreamingReceiver, + receiveMode: InternalReceiveMode.peekLock, + lockRenewer: undefined }); assert.isTrue(spy.called, "We do still call init() on the receiver"); @@ -270,7 +292,8 @@ describe("StreamingReceiver unit tests", () => { const receiverImpl = new ServiceBusReceiverImpl( createConnectionContextForTests(), "fakeEntityPath", - "peekLock" + "peekLock", + 1 ); closeables.push(receiverImpl); @@ -333,7 +356,9 @@ describe("StreamingReceiver unit tests", () => { } } as any) as StreamingReceiver; }, - abortSignal: abortController.signal + abortSignal: abortController.signal, + lockRenewer: undefined, + receiveMode: InternalReceiveMode.receiveAndDelete }); assert.isTrue(wasCalled); diff --git a/sdk/servicebus/service-bus/test/renewLock.spec.ts b/sdk/servicebus/service-bus/test/renewLock.spec.ts index dacc1f791149..c8ae23944ff4 100644 --- a/sdk/servicebus/service-bus/test/renewLock.spec.ts +++ b/sdk/servicebus/service-bus/test/renewLock.spec.ts @@ -10,7 +10,8 @@ import { TestMessage } from "./utils/testUtils"; import { ServiceBusClientForTests, createServiceBusClientForTests, - getRandomTestClientTypeWithNoSessions + getRandomTestClientTypeWithNoSessions, + AutoGeneratedEntity } from "./utils/testutils2"; import { ServiceBusReceiver } from "../src/receivers/receiver"; import { ServiceBusSender } from "../src/sender"; @@ -22,6 +23,7 @@ describe("Message Lock Renewal", () => { let receiver: ServiceBusReceiver; const testClientType = getRandomTestClientTypeWithNoSessions(); + let autoGeneratedEntity: AutoGeneratedEntity; before(() => { serviceBusClient = createServiceBusClientForTests(); @@ -32,11 +34,11 @@ describe("Message Lock Renewal", () => { }); beforeEach(async () => { - const entityNames = await serviceBusClient.test.createTestEntities(testClientType); - receiver = await serviceBusClient.test.createPeekLockReceiver(entityNames); + autoGeneratedEntity = await serviceBusClient.test.createTestEntities(testClientType); + receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity); sender = serviceBusClient.test.addToCleanup( - serviceBusClient.createSender(entityNames.queue ?? entityNames.topic!) + serviceBusClient.createSender(autoGeneratedEntity.queue ?? autoGeneratedEntity.topic!) ); }); @@ -65,51 +67,61 @@ describe("Message Lock Renewal", () => { } ); - it( - testClientType + - ": Streaming Receiver: complete() after lock expiry with auto-renewal disabled throws error", - async function(): Promise { - await testAutoLockRenewalConfigBehavior(sender, receiver, { - maxAutoRenewDurationInMs: 0, - delayBeforeAttemptingToCompleteMessageInSeconds: 31, - willCompleteFail: true + const receiveMethodType: ("subscribe" | "receive" | "iterator")[] = [ + "iterator", + "subscribe", + "receive" + ]; + + describe(`Using configurable renew durations`, () => { + receiveMethodType.forEach((receiveMethodType) => { + it(`${testClientType}: [${receiveMethodType}] Streaming Receiver: complete() after lock expiry with auto-renewal disabled throws error`, async function(): Promise< + void + > { + await testAutoLockRenewalConfigBehavior(sender, receiveMethodType, { + maxAutoRenewDurationInMs: 0, + delayBeforeAttemptingToCompleteMessageInSeconds: 31, + willCompleteFail: true + }); }); - } - ); + }); - it( - testClientType + ": Streaming Receiver: lock will not expire until configured time", - async function(): Promise { - await testAutoLockRenewalConfigBehavior(sender, receiver, { - maxAutoRenewDurationInMs: 38 * 1000, - delayBeforeAttemptingToCompleteMessageInSeconds: 35, - willCompleteFail: false + receiveMethodType.forEach((receiveMethodType) => { + it(`${testClientType}: [${receiveMethodType}] : Streaming Receiver: lock will not expire until configured time`, async function(): Promise< + void + > { + await testAutoLockRenewalConfigBehavior(sender, receiveMethodType, { + maxAutoRenewDurationInMs: 38 * 1000, + delayBeforeAttemptingToCompleteMessageInSeconds: 35, + willCompleteFail: false + }); }); - } - ); + }); - it( - testClientType + ": Streaming Receiver: lock expires sometime after configured time", - async function(): Promise { - await testAutoLockRenewalConfigBehavior(sender, receiver, { - maxAutoRenewDurationInMs: 35 * 1000, - delayBeforeAttemptingToCompleteMessageInSeconds: 55, - willCompleteFail: true - }); - } - ).timeout(95000 + 30000); + receiveMethodType.forEach((receiveMethodType) => { + it(`${testClientType}: [${receiveMethodType}] : Streaming Receiver: lock expires sometime after configured time`, async function(): Promise< + void + > { + await testAutoLockRenewalConfigBehavior(sender, receiveMethodType, { + maxAutoRenewDurationInMs: 35 * 1000, + delayBeforeAttemptingToCompleteMessageInSeconds: 55, + willCompleteFail: true + }); + }).timeout(95000 + 30000); + }); - it( - testClientType + - ": Streaming Receiver: No lock renewal when config value is less than lock duration", - async function(): Promise { - await testAutoLockRenewalConfigBehavior(sender, receiver, { - maxAutoRenewDurationInMs: 15 * 1000, - delayBeforeAttemptingToCompleteMessageInSeconds: 31, - willCompleteFail: true + receiveMethodType.forEach((receiveMethodType) => { + it(`${testClientType}: [${receiveMethodType}] Streaming Receiver: No lock renewal when config value is less than lock duration`, async function(): Promise< + void + > { + await testAutoLockRenewalConfigBehavior(sender, receiveMethodType, { + maxAutoRenewDurationInMs: 15 * 1000, + delayBeforeAttemptingToCompleteMessageInSeconds: 31, + willCompleteFail: true + }); }); - } - ); + }); + }); const lockDurationInMilliseconds = 30000; @@ -261,8 +273,7 @@ describe("Message Lock Renewal", () => { receiver.subscribe( { processMessage, processError }, { - autoComplete: false, - maxAutoRenewLockDurationInMs: 0 + autoComplete: false } ); await delay(10000); @@ -283,70 +294,90 @@ describe("Message Lock Renewal", () => { async function testAutoLockRenewalConfigBehavior( sender: ServiceBusSender, - receiver: ServiceBusReceiver, + type: "subscribe" | "receive" | "iterator", options: AutoLockRenewalTestOptions ): Promise { - let numOfMessagesReceived = 0; - const testMessage = TestMessage.getSample(); - await sender.sendMessages(testMessage); + const expectedMessage = TestMessage.getSample(`${type} ${Date.now().toString()}`); + await sender.sendMessages(expectedMessage); - async function processMessage( - brokeredMessage: ServiceBusReceivedMessageWithLock - ): Promise { - if (numOfMessagesReceived < 1) { - numOfMessagesReceived++; + await receiver.close(); - should.equal( - brokeredMessage.body, - testMessage.body, - "MessageBody is different than expected" - ); - should.equal( - brokeredMessage.messageId, - testMessage.messageId, - "MessageId is different than expected" - ); + receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { + maxLockAutoRenewDurationInMs: options.maxAutoRenewDurationInMs + }); - // Sleeping... - await delay(options.delayBeforeAttemptingToCompleteMessageInSeconds * 1000); + try { + const actualMessage = await receiveSingleMessageUsingSpecificReceiveMethod(type); - let errorWasThrown: boolean = false; - await brokeredMessage.complete().catch((err) => { - should.equal(err.code, "MessageLockLostError", "Error code is different than expected"); - errorWasThrown = true; - }); + should.equal( + actualMessage.body, + expectedMessage.body, + "MessageBody is different than expected" + ); + should.equal( + actualMessage.messageId, + expectedMessage.messageId, + "MessageId is different than expected" + ); - should.equal(errorWasThrown, options.willCompleteFail, "Error Thrown flag value mismatch"); - } - } + // Sleeping... + await delay(options.delayBeforeAttemptingToCompleteMessageInSeconds * 1000); - receiver.subscribe( - { processMessage, processError }, - { - autoComplete: false, - maxAutoRenewLockDurationInMs: options.maxAutoRenewDurationInMs - } - ); - await delay(options.delayBeforeAttemptingToCompleteMessageInSeconds * 1000 + 10000); + try { + await actualMessage.complete(); - if (uncaughtErrorFromHandlers) { - chai.assert.fail(uncaughtErrorFromHandlers.message); + if (options.willCompleteFail) { + should.fail("complete() should throw an error"); + } + } catch (err) { + if (options.willCompleteFail) { + should.equal(err.code, "MessageLockLostError", "Error code is different than expected"); + } else { + throw err; + } + } + } finally { + try { + const purgingReceiver = await serviceBusClient.test.createReceiveAndDeleteReceiver( + autoGeneratedEntity + ); + await purgingReceiver.receiveMessages(10, { maxWaitTimeInMs: 1000 }); + } catch (err) { + // ignore these errors + } } + } - should.equal(numOfMessagesReceived, 1, "Mismatch in number of messages received"); - - if (options.willCompleteFail) { - // Clean up any left over messages - await receiver.close(); - - receiver = await serviceBusClient.test.createPeekLockReceiver( - await serviceBusClient.test.createTestEntities(testClientType) - ); - - const unprocessedMsgsBatch = await receiver.receiveMessages(1); - if (unprocessedMsgsBatch.length) { - await unprocessedMsgsBatch[0].complete(); + async function receiveSingleMessageUsingSpecificReceiveMethod( + type: "subscribe" | "receive" | "iterator" + ): Promise { + switch (type) { + case "subscribe": { + return await new Promise((resolve, reject) => { + receiver.subscribe( + { + processMessage: async (msg) => resolve(msg), + processError: async (err) => reject(err) + }, + { + autoComplete: false + } + ); + }); + } + case "receive": { + const messages = await receiver.receiveMessages(1); + should.equal(1, messages.length); + return messages[0]; + } + case "iterator": { + for await (const message of receiver.getMessageIterator()) { + return message; + } + throw new Error("Failed to get message using iterator"); } + default: + throw new Error(`No receive method corresponds to type ${type}`); } } diff --git a/sdk/servicebus/service-bus/test/retries.spec.ts b/sdk/servicebus/service-bus/test/retries.spec.ts index 6df6f6823b51..e0fed6170b3d 100644 --- a/sdk/servicebus/service-bus/test/retries.spec.ts +++ b/sdk/servicebus/service-bus/test/retries.spec.ts @@ -18,6 +18,7 @@ import { ServiceBusSessionReceiver } from "../src/receivers/sessionReceiver"; import { ServiceBusReceiver, ServiceBusReceiverImpl } from "../src/receivers/receiver"; +import { InternalReceiveMode } from "../src/serviceBusMessage"; describe("Retries - ManagementClient", () => { let sender: ServiceBusSender; @@ -353,7 +354,11 @@ describe("Retries - Receive methods", () => { // Mocking batchingReceiver.receive to throw the error and fail const batchingReceiver = BatchingReceiver.create( (receiver as any)._context, - "dummyEntityPath" + "dummyEntityPath", + { + lockRenewer: undefined, + receiveMode: InternalReceiveMode.peekLock + } ); batchingReceiver.isOpen = () => true; batchingReceiver.receive = fakeFunction; diff --git a/sdk/servicebus/service-bus/test/streamingReceiver.spec.ts b/sdk/servicebus/service-bus/test/streamingReceiver.spec.ts index b03f360f466e..e2a6da6d1713 100644 --- a/sdk/servicebus/service-bus/test/streamingReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/streamingReceiver.spec.ts @@ -171,7 +171,8 @@ describe("Streaming Receiver Tests", () => { (receiver as any)._context, receiver.entityPath, { - receiveMode: InternalReceiveMode.peekLock + receiveMode: InternalReceiveMode.peekLock, + lockRenewer: undefined } ); diff --git a/sdk/servicebus/service-bus/test/utils/testUtils.ts b/sdk/servicebus/service-bus/test/utils/testUtils.ts index 50303c661d9c..b7f3e3092d08 100644 --- a/sdk/servicebus/service-bus/test/utils/testUtils.ts +++ b/sdk/servicebus/service-bus/test/utils/testUtils.ts @@ -9,18 +9,21 @@ dotenv.config(); export class TestMessage { static sessionId: string = "my-session"; - static getSample(): ServiceBusMessage { - const randomNumber = Math.random(); + static getSample(randomTag?: string): ServiceBusMessage { + if (randomTag == null) { + randomTag = Math.random().toString(); + } + return { - body: `message body ${randomNumber}`, - messageId: `message id ${randomNumber}`, + body: `message body ${randomTag}`, + messageId: `message id ${randomTag}`, partitionKey: `dummy partition key`, - contentType: `content type ${randomNumber}`, - correlationId: `correlation id ${randomNumber}`, + contentType: `content type ${randomTag}`, + correlationId: `correlation id ${randomTag}`, timeToLive: 60 * 60 * 24, - label: `label ${randomNumber}`, - to: `to ${randomNumber}`, - replyTo: `reply to ${randomNumber}`, + label: `label ${randomTag}`, + to: `to ${randomTag}`, + replyTo: `reply to ${randomTag}`, scheduledEnqueueTimeUtc: new Date(), properties: { propOne: 1, diff --git a/sdk/servicebus/service-bus/test/utils/testutils2.ts b/sdk/servicebus/service-bus/test/utils/testutils2.ts index d1e2276d184d..4ab7146fefab 100644 --- a/sdk/servicebus/service-bus/test/utils/testutils2.ts +++ b/sdk/servicebus/service-bus/test/utils/testutils2.ts @@ -8,7 +8,8 @@ import { ServiceBusReceiver, ServiceBusSessionReceiver, ServiceBusClientOptions, - AcceptSessionOptions + AcceptSessionOptions, + CreateReceiverOptions } from "../../src"; import { TestClientType, TestMessage } from "./testUtils"; @@ -31,17 +32,24 @@ dotenv.config(); const env = getEnvVars(); const should = chai.should(); -const defaultLockDuration = "PT30S"; // 30 seconds in ISO 8601 FORMAT - equivalent to "P0Y0M0DT0H0M30S" - -function getEntityNames( - testClientType: TestClientType -): { +/** + * Identifier of an auto-generated entity. + * + * note: For the entity name itself either 'queue' or 'topic' (or 'topic and 'subscription') + * will be filled out, depending on the type of test entity you created. + * + */ +export interface AutoGeneratedEntity { queue?: string; topic?: string; subscription?: string; usesSessions: boolean; isPartitioned: boolean; -} { +} + +const defaultLockDuration = "PT30S"; // 30 seconds in ISO 8601 FORMAT - equivalent to "P0Y0M0DT0H0M30S" + +function getEntityNames(testClientType: TestClientType): AutoGeneratedEntity { const name = testClientType; let prefix = ""; let isPartitioned = false; @@ -322,24 +330,25 @@ export class ServiceBusTestHelpers { * The receiver created by this method will be cleaned up by `afterEach()` */ async createPeekLockReceiver( - entityNames: Omit, "isPartitioned"> + entityNames: Omit, "isPartitioned">, + options?: CreateReceiverOptions<"peekLock"> ): Promise> { - try { + if (entityNames.usesSessions) { // if you're creating a receiver this way then you'll just use the default // session ID for your receiver. // if you want to get more specific use the `getPeekLockSessionReceiver` method // instead. - return await this.acceptSessionWithPeekLock(entityNames, TestMessage.sessionId); - } catch (err) { - if (!(err instanceof TypeError)) { - throw err; - } + return await this.acceptSessionWithPeekLock(entityNames, TestMessage.sessionId, options); } return this.addToCleanup( entityNames.queue - ? this._serviceBusClient.createReceiver(entityNames.queue) - : this._serviceBusClient.createReceiver(entityNames.topic!, entityNames.subscription!) + ? this._serviceBusClient.createReceiver(entityNames.queue, options) + : this._serviceBusClient.createReceiver( + entityNames.topic!, + entityNames.subscription!, + options + ) ); } From 766495a78b6fffcc0cfabc6607da7ebd5785cc92 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Mon, 5 Oct 2020 18:03:47 -0700 Subject: [PATCH 02/14] Chicken and the egg --- sdk/servicebus/service-bus/CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/service-bus/CHANGELOG.md b/sdk/servicebus/service-bus/CHANGELOG.md index 010e87aa4bca..2a0e6ef472c5 100644 --- a/sdk/servicebus/service-bus/CHANGELOG.md +++ b/sdk/servicebus/service-bus/CHANGELOG.md @@ -12,7 +12,7 @@ [PR 11117](https://github.com/Azure/azure-sdk-for-js/pull/11117) - Message locks can be auto-renewed in all receive methods (receiver.receiveMessages, receiver.subcribe and receiver.getMessageIterator). This can be configured in options when calling `ServiceBusClient.createReceiver()`. - [PR]() + [PR 11658](https://github.com/Azure/azure-sdk-for-js/pull/11658) ### Breaking changes From ef3b32443671c13c5c632b534fee5f95faf3c3a0 Mon Sep 17 00:00:00 2001 From: Ramya Achutha Rao Date: Mon, 5 Oct 2020 22:26:45 -0700 Subject: [PATCH 03/14] Remove unused import --- .../service-bus/test/internal/batchingReceiver.spec.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts b/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts index df8c24aa2701..79e4dc669bf0 100644 --- a/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts @@ -30,7 +30,6 @@ import { StandardAbortMessage } from "../../src/util/utils"; import { OnAmqpEventAsPromise } from "../../src/core/messageReceiver"; import { ConnectionContext } from "../../src/connectionContext"; import { ServiceBusReceiverImpl } from "../../src/receivers/receiver"; -import { LockRenewer } from "../../src/core/autoLockRenewer"; describe("BatchingReceiver unit tests", () => { let closeables: { close(): Promise }[]; From 547e879e2f47a221f19b3c9265e3a821894199f8 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 10:11:27 -0700 Subject: [PATCH 04/14] Disable auto-lock renewal for batchreceiver tests. It interferes with at least one test and for compatibility having it off is simpler. We can revisit later - there is coverage for this case in renewLock.spec.ts so it's not untested. --- sdk/servicebus/service-bus/test/batchReceiver.spec.ts | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts index 638293d3daed..93ca44695801 100644 --- a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts @@ -40,7 +40,12 @@ let deadLetterReceiver: ServiceBusReceiver; async function beforeEachTest(entityType: TestClientType): Promise { entityNames = await serviceBusClient.test.createTestEntities(entityType); - receiver = await serviceBusClient.test.createPeekLockReceiver(entityNames); + receiver = await serviceBusClient.test.createPeekLockReceiver(entityNames, { + // prior to a recent change the behavior was always to _not_ auto-renew locks. + // for compat with these tests I'm just disabling this. There are tests in renewLocks.spec.ts that + // ensure lock renewal does work with batching. + maxLockAutoRenewDurationInMs: 0 + }); sender = serviceBusClient.test.addToCleanup( serviceBusClient.createSender(entityNames.queue ?? entityNames.topic!) @@ -777,7 +782,7 @@ describe("Batching Receiver", () => { await batch[0].complete(); } - it( + it.only( noSessionTestClientType + ": No settlement of the message is retained with incremented deliveryCount", async function(): Promise { From ac45a8737f65de03399e6541c91c128de2518946 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 10:12:27 -0700 Subject: [PATCH 05/14] Remove .only --- sdk/servicebus/service-bus/test/batchReceiver.spec.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts index 93ca44695801..6ab25e54a5c0 100644 --- a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts @@ -782,7 +782,7 @@ describe("Batching Receiver", () => { await batch[0].complete(); } - it.only( + it( noSessionTestClientType + ": No settlement of the message is retained with incremented deliveryCount", async function(): Promise { From bd1409929eac3f3a91f4f6d654665fc72c39c498 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 10:12:52 -0700 Subject: [PATCH 06/14] Remove unused import --- .../service-bus/test/internal/batchingReceiver.spec.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts b/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts index df8c24aa2701..79e4dc669bf0 100644 --- a/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/internal/batchingReceiver.spec.ts @@ -30,7 +30,6 @@ import { StandardAbortMessage } from "../../src/util/utils"; import { OnAmqpEventAsPromise } from "../../src/core/messageReceiver"; import { ConnectionContext } from "../../src/connectionContext"; import { ServiceBusReceiverImpl } from "../../src/receivers/receiver"; -import { LockRenewer } from "../../src/core/autoLockRenewer"; describe("BatchingReceiver unit tests", () => { let closeables: { close(): Promise }[]; From ba136668dc4badc5eeabb685737dce5b9a315f5d Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 11:42:14 -0700 Subject: [PATCH 07/14] Updating the renew lock tests. There were a few tests that were relying on the fact that batching didn't use to renew locks (it does now). Changed so each test creates their own receiver, rather than re-using a single one for all tests. --- .../service-bus/test/renewLock.spec.ts | 39 +++++++++++-------- 1 file changed, 23 insertions(+), 16 deletions(-) diff --git a/sdk/servicebus/service-bus/test/renewLock.spec.ts b/sdk/servicebus/service-bus/test/renewLock.spec.ts index c8ae23944ff4..36815e256d28 100644 --- a/sdk/servicebus/service-bus/test/renewLock.spec.ts +++ b/sdk/servicebus/service-bus/test/renewLock.spec.ts @@ -20,7 +20,6 @@ import { ServiceBusReceivedMessageWithLock } from "../src/serviceBusMessage"; describe("Message Lock Renewal", () => { let serviceBusClient: ServiceBusClientForTests; let sender: ServiceBusSender; - let receiver: ServiceBusReceiver; const testClientType = getRandomTestClientTypeWithNoSessions(); let autoGeneratedEntity: AutoGeneratedEntity; @@ -35,7 +34,6 @@ describe("Message Lock Renewal", () => { beforeEach(async () => { autoGeneratedEntity = await serviceBusClient.test.createTestEntities(testClientType); - receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity); sender = serviceBusClient.test.addToCleanup( serviceBusClient.createSender(autoGeneratedEntity.queue ?? autoGeneratedEntity.topic!) @@ -49,21 +47,21 @@ describe("Message Lock Renewal", () => { it( testClientType + ": Batch Receiver: renewLock() resets lock duration each time.", async function(): Promise { - await testBatchReceiverManualLockRenewalHappyCase(sender, receiver); + await testBatchReceiverManualLockRenewalHappyCase(sender); } ); it( testClientType + ": Batch Receiver: complete() after lock expiry with throws error", async function(): Promise { - await testBatchReceiverManualLockRenewalErrorOnLockExpiry(sender, receiver); + await testBatchReceiverManualLockRenewalErrorOnLockExpiry(sender); } ); it( testClientType + ": Streaming Receiver: renewLock() resets lock duration each time.", async function(): Promise { - await testStreamingReceiverManualLockRenewalHappyCase(sender, receiver); + await testStreamingReceiverManualLockRenewalHappyCase(sender); } ); @@ -135,9 +133,12 @@ describe("Message Lock Renewal", () => { * Test renewLock() after receiving a message using Batch Receiver */ async function testBatchReceiverManualLockRenewalHappyCase( - sender: ServiceBusSender, - receiver: ServiceBusReceiver + sender: ServiceBusSender ): Promise { + const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { + maxLockAutoRenewDurationInMs: 0 + }); + const testMessage = TestMessage.getSample(); await sender.sendMessages(testMessage); @@ -183,9 +184,12 @@ describe("Message Lock Renewal", () => { * Test settling of message from Batch Receiver fails after message lock expires */ async function testBatchReceiverManualLockRenewalErrorOnLockExpiry( - sender: ServiceBusSender, - receiver: ServiceBusReceiver + sender: ServiceBusSender ): Promise { + const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { + maxLockAutoRenewDurationInMs: 0 + }); + const testMessage = TestMessage.getSample(); await sender.sendMessages(testMessage); @@ -196,7 +200,8 @@ describe("Message Lock Renewal", () => { should.equal(msgs[0].body, testMessage.body, "MessageBody is different than expected"); should.equal(msgs[0].messageId, testMessage.messageId, "MessageId is different than expected"); - // Sleeping 30 seconds... + // sleeping for long enough to let the message expire (we assume + // the remote entity is using the default duration of 30 seconds for their locks) await delay(lockDurationInMilliseconds + 1000); let errorWasThrown: boolean = false; @@ -216,9 +221,12 @@ describe("Message Lock Renewal", () => { * Test renewLock() after receiving a message using Streaming Receiver with autoLockRenewal disabled */ async function testStreamingReceiverManualLockRenewalHappyCase( - sender: ServiceBusSender, - receiver: ServiceBusReceiver + sender: ServiceBusSender ): Promise { + const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { + maxLockAutoRenewDurationInMs: 0 + }); + let numOfMessagesReceived = 0; const testMessage = TestMessage.getSample(); await sender.sendMessages(testMessage); @@ -300,14 +308,12 @@ describe("Message Lock Renewal", () => { const expectedMessage = TestMessage.getSample(`${type} ${Date.now().toString()}`); await sender.sendMessages(expectedMessage); - await receiver.close(); - - receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { + const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { maxLockAutoRenewDurationInMs: options.maxAutoRenewDurationInMs }); try { - const actualMessage = await receiveSingleMessageUsingSpecificReceiveMethod(type); + const actualMessage = await receiveSingleMessageUsingSpecificReceiveMethod(receiver, type); should.equal( actualMessage.body, @@ -349,6 +355,7 @@ describe("Message Lock Renewal", () => { } async function receiveSingleMessageUsingSpecificReceiveMethod( + receiver: ServiceBusReceiver, type: "subscribe" | "receive" | "iterator" ): Promise { switch (type) { From cf328d5bbfdc6264f497bb2adce972aebf85665c Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 11:47:20 -0700 Subject: [PATCH 08/14] Rename the message lock renewal duration option to be more specific to messages. --- sdk/servicebus/service-bus/src/models.ts | 2 +- sdk/servicebus/service-bus/src/serviceBusClient.ts | 4 ++-- sdk/servicebus/service-bus/test/batchReceiver.spec.ts | 2 +- sdk/servicebus/service-bus/test/renewLock.spec.ts | 8 ++++---- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/sdk/servicebus/service-bus/src/models.ts b/sdk/servicebus/service-bus/src/models.ts index 06cb8cc34a57..89bc10030cc0 100644 --- a/sdk/servicebus/service-bus/src/models.ts +++ b/sdk/servicebus/service-bus/src/models.ts @@ -89,7 +89,7 @@ export interface CreateReceiverOptions { * - **Default**: `300 * 1000` milliseconds (5 minutes). * - **To disable autolock renewal**, set this to `0`. */ - maxLockAutoRenewDurationInMs?: number; + maxMessageLockRenewalDurationInMs?: number; } /** diff --git a/sdk/servicebus/service-bus/src/serviceBusClient.ts b/sdk/servicebus/service-bus/src/serviceBusClient.ts index e6103344331d..1bd4240b955a 100644 --- a/sdk/servicebus/service-bus/src/serviceBusClient.ts +++ b/sdk/servicebus/service-bus/src/serviceBusClient.ts @@ -237,8 +237,8 @@ export class ServiceBusClient { } const maxLockAutoRenewDurationInMs = - options?.maxLockAutoRenewDurationInMs != null - ? options.maxLockAutoRenewDurationInMs + options?.maxMessageLockRenewalDurationInMs != null + ? options.maxMessageLockRenewalDurationInMs : 5 * 60 * 1000; if (receiveMode === "peekLock") { diff --git a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts index 6ab25e54a5c0..725bd28ed213 100644 --- a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts @@ -44,7 +44,7 @@ async function beforeEachTest(entityType: TestClientType): Promise { // prior to a recent change the behavior was always to _not_ auto-renew locks. // for compat with these tests I'm just disabling this. There are tests in renewLocks.spec.ts that // ensure lock renewal does work with batching. - maxLockAutoRenewDurationInMs: 0 + maxMessageLockRenewalDurationInMs: 0 }); sender = serviceBusClient.test.addToCleanup( diff --git a/sdk/servicebus/service-bus/test/renewLock.spec.ts b/sdk/servicebus/service-bus/test/renewLock.spec.ts index 36815e256d28..bab7c96722d8 100644 --- a/sdk/servicebus/service-bus/test/renewLock.spec.ts +++ b/sdk/servicebus/service-bus/test/renewLock.spec.ts @@ -136,7 +136,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxLockAutoRenewDurationInMs: 0 + maxMessageLockRenewalDurationInMs: 0 }); const testMessage = TestMessage.getSample(); @@ -187,7 +187,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxLockAutoRenewDurationInMs: 0 + maxMessageLockRenewalDurationInMs: 0 }); const testMessage = TestMessage.getSample(); @@ -224,7 +224,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxLockAutoRenewDurationInMs: 0 + maxMessageLockRenewalDurationInMs: 0 }); let numOfMessagesReceived = 0; @@ -309,7 +309,7 @@ describe("Message Lock Renewal", () => { await sender.sendMessages(expectedMessage); const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxLockAutoRenewDurationInMs: options.maxAutoRenewDurationInMs + maxMessageLockRenewalDurationInMs: options.maxAutoRenewDurationInMs }); try { From b93ec03ac53329e9661459dfaecf7d5af958feb0 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 11:51:34 -0700 Subject: [PATCH 09/14] Remove the logError call in batchingReceiver when the renewer fails. The renewer already has a great log message in place. --- sdk/servicebus/service-bus/src/core/batchingReceiver.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/sdk/servicebus/service-bus/src/core/batchingReceiver.ts b/sdk/servicebus/service-bus/src/core/batchingReceiver.ts index aae238925309..4d0c17db8b69 100644 --- a/sdk/servicebus/service-bus/src/core/batchingReceiver.ts +++ b/sdk/servicebus/service-bus/src/core/batchingReceiver.ts @@ -127,8 +127,9 @@ export class BatchingReceiver extends MessageReceiver { if (this._lockRenewer) { for (const message of messages) { - this._lockRenewer.start(this, message, (error) => { - logError(error, `${this.logPrefix} Failed to renew lock for message.`); + this._lockRenewer.start(this, message, (_error) => { + // the auto lock renewer already logs this in a detailed way. So this hook is mainly here + // to potentially forward the error to the user (which we're not doing yet) }); } } From da812c60789ad1bce654693e2552c17df4b09e71 Mon Sep 17 00:00:00 2001 From: Richard Park <51494936+richardpark-msft@users.noreply.github.com> Date: Tue, 6 Oct 2020 11:53:59 -0700 Subject: [PATCH 10/14] It is per link indeed. :) Co-authored-by: Ramya Rao --- sdk/servicebus/service-bus/src/core/autoLockRenewer.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts index 9e04a39e0156..8fc13897450a 100644 --- a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts +++ b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts @@ -73,7 +73,7 @@ export class LockRenewer { } /** - * Cancels all pending lock renewals and removes all entries from our internal cache. + * Cancels all pending lock renewals for messages on given link and removes all entries from our internal cache. */ stopAll(linkEntity: MinimalLink) { logger.verbose( From e8fc11bdbd305727e27c875c35a1895ba9c7acf3 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 12:41:14 -0700 Subject: [PATCH 11/14] setting name has changed. --- sdk/servicebus/service-bus/src/core/autoLockRenewer.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts index 9e04a39e0156..817eec21c5b3 100644 --- a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts +++ b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts @@ -54,7 +54,7 @@ export class LockRenewer { * @param context The connection context for your link entity (probably 'this._context') * @param options The ReceiveOptions passed through to your message receiver. * @returns if the lock mode is peek lock (or if is unspecified, thus defaulting to peekLock) - * and the options.maxAutoRenewLockDurationInMs is > 0..Otherwise, returns undefined. + * and the options.maxMessageLockRenewalDurationInMs is > 0..Otherwise, returns undefined. */ static create( context: Pick, From c569f732aa6dd4dbb418859921d18b670cacdb77 Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 12:47:24 -0700 Subject: [PATCH 12/14] Updating name in all sorts of places for consistency with other SDKs. --- sdk/servicebus/service-bus/src/core/autoLockRenewer.ts | 2 +- sdk/servicebus/service-bus/src/models.ts | 2 +- sdk/servicebus/service-bus/src/serviceBusClient.ts | 4 ++-- sdk/servicebus/service-bus/test/batchReceiver.spec.ts | 2 +- sdk/servicebus/service-bus/test/renewLock.spec.ts | 8 ++++---- 5 files changed, 9 insertions(+), 9 deletions(-) diff --git a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts index e09356ce2163..dd1beb958e05 100644 --- a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts +++ b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts @@ -54,7 +54,7 @@ export class LockRenewer { * @param context The connection context for your link entity (probably 'this._context') * @param options The ReceiveOptions passed through to your message receiver. * @returns if the lock mode is peek lock (or if is unspecified, thus defaulting to peekLock) - * and the options.maxMessageLockRenewalDurationInMs is > 0..Otherwise, returns undefined. + * and the options.maxMessageLockAutoRenewDurationInMs is > 0..Otherwise, returns undefined. */ static create( context: Pick, diff --git a/sdk/servicebus/service-bus/src/models.ts b/sdk/servicebus/service-bus/src/models.ts index 89bc10030cc0..51207d16eb6e 100644 --- a/sdk/servicebus/service-bus/src/models.ts +++ b/sdk/servicebus/service-bus/src/models.ts @@ -89,7 +89,7 @@ export interface CreateReceiverOptions { * - **Default**: `300 * 1000` milliseconds (5 minutes). * - **To disable autolock renewal**, set this to `0`. */ - maxMessageLockRenewalDurationInMs?: number; + maxMessageLockAutoRenewDurationInMs?: number; } /** diff --git a/sdk/servicebus/service-bus/src/serviceBusClient.ts b/sdk/servicebus/service-bus/src/serviceBusClient.ts index 1bd4240b955a..900583c2b532 100644 --- a/sdk/servicebus/service-bus/src/serviceBusClient.ts +++ b/sdk/servicebus/service-bus/src/serviceBusClient.ts @@ -237,8 +237,8 @@ export class ServiceBusClient { } const maxLockAutoRenewDurationInMs = - options?.maxMessageLockRenewalDurationInMs != null - ? options.maxMessageLockRenewalDurationInMs + options?.maxMessageLockAutoRenewDurationInMs != null + ? options.maxMessageLockAutoRenewDurationInMs : 5 * 60 * 1000; if (receiveMode === "peekLock") { diff --git a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts index 725bd28ed213..6b2733d19eea 100644 --- a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts @@ -44,7 +44,7 @@ async function beforeEachTest(entityType: TestClientType): Promise { // prior to a recent change the behavior was always to _not_ auto-renew locks. // for compat with these tests I'm just disabling this. There are tests in renewLocks.spec.ts that // ensure lock renewal does work with batching. - maxMessageLockRenewalDurationInMs: 0 + maxMessageLockAutoRenewDurationInMs: 0 }); sender = serviceBusClient.test.addToCleanup( diff --git a/sdk/servicebus/service-bus/test/renewLock.spec.ts b/sdk/servicebus/service-bus/test/renewLock.spec.ts index bab7c96722d8..f7d9c7842cd4 100644 --- a/sdk/servicebus/service-bus/test/renewLock.spec.ts +++ b/sdk/servicebus/service-bus/test/renewLock.spec.ts @@ -136,7 +136,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockRenewalDurationInMs: 0 + maxMessageLockAutoRenewDurationInMs: 0 }); const testMessage = TestMessage.getSample(); @@ -187,7 +187,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockRenewalDurationInMs: 0 + maxMessageLockAutoRenewDurationInMs: 0 }); const testMessage = TestMessage.getSample(); @@ -224,7 +224,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockRenewalDurationInMs: 0 + maxMessageLockAutoRenewDurationInMs: 0 }); let numOfMessagesReceived = 0; @@ -309,7 +309,7 @@ describe("Message Lock Renewal", () => { await sender.sendMessages(expectedMessage); const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockRenewalDurationInMs: options.maxAutoRenewDurationInMs + maxMessageLockAutoRenewDurationInMs: options.maxAutoRenewDurationInMs }); try { From 9c23a72e52a5bdee51280aa4922f10387759aeac Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 12:55:22 -0700 Subject: [PATCH 13/14] renaming field for consistency with .net --- sdk/servicebus/service-bus/src/core/autoLockRenewer.ts | 2 +- sdk/servicebus/service-bus/src/models.ts | 2 +- sdk/servicebus/service-bus/src/serviceBusClient.ts | 4 ++-- sdk/servicebus/service-bus/test/batchReceiver.spec.ts | 2 +- sdk/servicebus/service-bus/test/renewLock.spec.ts | 8 ++++---- 5 files changed, 9 insertions(+), 9 deletions(-) diff --git a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts index dd1beb958e05..89ed7e30f985 100644 --- a/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts +++ b/sdk/servicebus/service-bus/src/core/autoLockRenewer.ts @@ -54,7 +54,7 @@ export class LockRenewer { * @param context The connection context for your link entity (probably 'this._context') * @param options The ReceiveOptions passed through to your message receiver. * @returns if the lock mode is peek lock (or if is unspecified, thus defaulting to peekLock) - * and the options.maxMessageLockAutoRenewDurationInMs is > 0..Otherwise, returns undefined. + * and the options.maxAutoLockRenewalDurationInMs is > 0..Otherwise, returns undefined. */ static create( context: Pick, diff --git a/sdk/servicebus/service-bus/src/models.ts b/sdk/servicebus/service-bus/src/models.ts index 51207d16eb6e..a1e5a6c1583d 100644 --- a/sdk/servicebus/service-bus/src/models.ts +++ b/sdk/servicebus/service-bus/src/models.ts @@ -89,7 +89,7 @@ export interface CreateReceiverOptions { * - **Default**: `300 * 1000` milliseconds (5 minutes). * - **To disable autolock renewal**, set this to `0`. */ - maxMessageLockAutoRenewDurationInMs?: number; + maxAutoLockRenewalDurationInMs?: number; } /** diff --git a/sdk/servicebus/service-bus/src/serviceBusClient.ts b/sdk/servicebus/service-bus/src/serviceBusClient.ts index 900583c2b532..1c158dd6dd67 100644 --- a/sdk/servicebus/service-bus/src/serviceBusClient.ts +++ b/sdk/servicebus/service-bus/src/serviceBusClient.ts @@ -237,8 +237,8 @@ export class ServiceBusClient { } const maxLockAutoRenewDurationInMs = - options?.maxMessageLockAutoRenewDurationInMs != null - ? options.maxMessageLockAutoRenewDurationInMs + options?.maxAutoLockRenewalDurationInMs != null + ? options.maxAutoLockRenewalDurationInMs : 5 * 60 * 1000; if (receiveMode === "peekLock") { diff --git a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts index 6b2733d19eea..157ebafaddcb 100644 --- a/sdk/servicebus/service-bus/test/batchReceiver.spec.ts +++ b/sdk/servicebus/service-bus/test/batchReceiver.spec.ts @@ -44,7 +44,7 @@ async function beforeEachTest(entityType: TestClientType): Promise { // prior to a recent change the behavior was always to _not_ auto-renew locks. // for compat with these tests I'm just disabling this. There are tests in renewLocks.spec.ts that // ensure lock renewal does work with batching. - maxMessageLockAutoRenewDurationInMs: 0 + maxAutoLockRenewalDurationInMs: 0 }); sender = serviceBusClient.test.addToCleanup( diff --git a/sdk/servicebus/service-bus/test/renewLock.spec.ts b/sdk/servicebus/service-bus/test/renewLock.spec.ts index f7d9c7842cd4..cbcd79ae517b 100644 --- a/sdk/servicebus/service-bus/test/renewLock.spec.ts +++ b/sdk/servicebus/service-bus/test/renewLock.spec.ts @@ -136,7 +136,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockAutoRenewDurationInMs: 0 + maxAutoLockRenewalDurationInMs: 0 }); const testMessage = TestMessage.getSample(); @@ -187,7 +187,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockAutoRenewDurationInMs: 0 + maxAutoLockRenewalDurationInMs: 0 }); const testMessage = TestMessage.getSample(); @@ -224,7 +224,7 @@ describe("Message Lock Renewal", () => { sender: ServiceBusSender ): Promise { const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockAutoRenewDurationInMs: 0 + maxAutoLockRenewalDurationInMs: 0 }); let numOfMessagesReceived = 0; @@ -309,7 +309,7 @@ describe("Message Lock Renewal", () => { await sender.sendMessages(expectedMessage); const receiver = await serviceBusClient.test.createPeekLockReceiver(autoGeneratedEntity, { - maxMessageLockAutoRenewDurationInMs: options.maxAutoRenewDurationInMs + maxAutoLockRenewalDurationInMs: options.maxAutoRenewDurationInMs }); try { From c2ba8efd75f4bcc222bd90d2f5ad7c59d776d41a Mon Sep 17 00:00:00 2001 From: Richard Park Date: Tue, 6 Oct 2020 13:06:17 -0700 Subject: [PATCH 14/14] API file has also changed --- sdk/servicebus/service-bus/review/service-bus.api.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/service-bus/review/service-bus.api.md b/sdk/servicebus/service-bus/review/service-bus.api.md index a9c3d83631ed..651db13a0042 100644 --- a/sdk/servicebus/service-bus/review/service-bus.api.md +++ b/sdk/servicebus/service-bus/review/service-bus.api.md @@ -126,7 +126,7 @@ export interface CreateQueueOptions extends OperationOptions { // @public export interface CreateReceiverOptions { - maxLockAutoRenewDurationInMs?: number; + maxAutoLockRenewalDurationInMs?: number; receiveMode?: ReceiveModeT; subQueue?: SubQueue; }