From b0479363910462050abaa8063f36ba45655e938a Mon Sep 17 00:00:00 2001 From: Hemant Tanwar Date: Thu, 22 Oct 2020 13:01:06 -0700 Subject: [PATCH 1/6] testing --- .../ServiceBusSenderAsyncClient.java | 56 ++++++++++++++++++- .../ServiceBusSenderAsyncClientTest.java | 36 ++++++++++++ 2 files changed, 90 insertions(+), 2 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java index 5824c288c2eb..edfb60d3b8a6 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java @@ -36,6 +36,7 @@ import java.util.Objects; import java.util.Set; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; import java.util.function.BinaryOperator; @@ -371,15 +372,66 @@ public Flux scheduleMessages(Iterable messages, OffsetD return fluxError(logger, new NullPointerException("'scheduledEnqueueTime' cannot be null.")); } - return createBatch().flatMapMany(messageBatch -> { - messages.forEach(message -> messageBatch.tryAdd(message)); + return createBatch().flatMapMany(messageBatch -> { + AtomicInteger index = new AtomicInteger(); + messages.forEach(message -> { + index.incrementAndGet(); + boolean added = messageBatch.tryAdd(message); + if (!added) { + System.out.println(getClass().getName() + " !!!! ERROR could not add index =" + index); + throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, "Could not fit in one batch message", + null); + } + } + ); return getSendLink().flatMapMany(link -> connectionProcessor .flatMap(connection -> connection.getManagementNode(entityName, entityType)) .flatMapMany(managementNode -> managementNode.schedule(messageBatch.getMessages(), scheduledEnqueueTime, messageBatch.getMaxSizeInBytes(), link.getLinkName(), transactionContext))); }); + + return getSendLink() + .flatMap(link -> link.getLinkSize() + .flatMap(size -> { + final int batchSize = size > 0 ? size : MAX_MESSAGE_LENGTH_BYTES; + return createBatch().map(messageBatch -> { + //AtomicInteger index = new AtomicInteger(); + messages.forEach(message -> { + //index.incrementAndGet(); + boolean added = messageBatch.tryAdd(message); + if (!added) { + System.out.println(getClass().getName() + " !!!! ERROR could not add index =" ); + throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, "Could not fit in one batch message", + null); + } + } + ); + return connectionProcessor + .flatMap(connection -> connection.getManagementNode(entityName, entityType)) + .flatMapMany(managementNode -> managementNode.schedule(messageBatch.getMessages(), scheduledEnqueueTime, + messageBatch.getMaxSizeInBytes(), link.getLinkName(), transactionContext)); + })); + + + /* final CreateBatchOptions batchOptions = new CreateBatchOptions() + .setMaximumSizeInBytes(batchSize); + return messages.collect(new AmqpMessageCollector(batchOptions, 1, + link::getErrorContext, tracerProvider, messageSerializer, entityName, + link.getHostname()));*/ + })); + } + /*{ + if (!messageBatch.tryAdd(message)) { + + final String errorMessage = String.format(Locale.US, + "EventData does not fit into maximum number of batches. '%s'", maxNumberOfBatches); + + throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, errorMessage, + contextProvider.getErrorContext()); + + };*/ /** * Cancels the enqueuing of an already scheduled message, if it was not already enqueued. * diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java index 4a15d56be719..ec9048938385 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java @@ -58,6 +58,7 @@ import java.util.Collections; import java.util.List; import java.util.UUID; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.IntStream; @@ -246,6 +247,41 @@ void createBatchWhenSizeTooBig() { .verify(); } + @Test + void scheduleMessageSizeTooBig() throws InterruptedException { + // Arrange + long sequenceNumberReturned = 10; + + int maxLinkSize = 1024; + int batchSize = maxLinkSize + 10; + + OffsetDateTime instant = mock(OffsetDateTime.class); + final List messages = TestUtils.getServiceBusMessages(batchSize, UUID.randomUUID().toString()); + + final AmqpSendLink link = mock(AmqpSendLink.class); + when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); + + when(connection.createSendLink(eq(ENTITY_NAME), eq(ENTITY_NAME), any(AmqpRetryOptions.class), isNull())) + .thenReturn(Mono.just(link)); + when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); + when(managementNode.schedule(anyList(), eq(instant), any(Integer.class), any(), isNull())) + .thenReturn(Flux.just(sequenceNumberReturned)); + + // Act & Assert + sender.scheduleMessages(messages, instant) + .subscribe(aLong -> { + System.out.println("!!!! Test Messages Scheduled : " + aLong); + }); + TimeUnit.SECONDS.sleep(5); + + verify(managementNode).schedule(sbMessagesCaptor.capture(), eq(instant), eq(MAX_MESSAGE_LENGTH_BYTES), eq(LINK_NAME), isNull()); + List actualMessages = sbMessagesCaptor.getValue(); + Assertions.assertNotNull(actualMessages); + System.out.println(" !!!! Test actualMessages : " + actualMessages.size()); + //Assertions.assertEquals(1, actualMessages.size()); + //Assertions.assertEquals(message, actualMessages.get(0)); + } + /** * Verifies that the producer can create a batch with a given {@link CreateBatchOptions#getMaximumSizeInBytes()}. */ From 154af1ae961188632c4240310b3550ac17712fb8 Mon Sep 17 00:00:00 2001 From: Hemant Tanwar Date: Thu, 22 Oct 2020 15:40:41 -0700 Subject: [PATCH 2/6] fix batch size --- .../ServiceBusSenderAsyncClient.java | 87 +++++++------------ .../ServiceBusSenderAsyncClientTest.java | 25 ++---- .../ServiceBusSessionManagerTest.java | 2 +- 3 files changed, 38 insertions(+), 76 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java index edfb60d3b8a6..f4b2f7f2959c 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java @@ -372,66 +372,37 @@ public Flux scheduleMessages(Iterable messages, OffsetD return fluxError(logger, new NullPointerException("'scheduledEnqueueTime' cannot be null.")); } - return createBatch().flatMapMany(messageBatch -> { - AtomicInteger index = new AtomicInteger(); - messages.forEach(message -> { - index.incrementAndGet(); - boolean added = messageBatch.tryAdd(message); - if (!added) { - System.out.println(getClass().getName() + " !!!! ERROR could not add index =" + index); - throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, "Could not fit in one batch message", - null); - } - } - ); - return getSendLink().flatMapMany(link -> connectionProcessor - .flatMap(connection -> connection.getManagementNode(entityName, entityType)) - .flatMapMany(managementNode -> managementNode.schedule(messageBatch.getMessages(), scheduledEnqueueTime, - messageBatch.getMaxSizeInBytes(), link.getLinkName(), transactionContext))); - }); - - return getSendLink() - .flatMap(link -> link.getLinkSize() - .flatMap(size -> { - final int batchSize = size > 0 ? size : MAX_MESSAGE_LENGTH_BYTES; - return createBatch().map(messageBatch -> { - //AtomicInteger index = new AtomicInteger(); - messages.forEach(message -> { - //index.incrementAndGet(); - boolean added = messageBatch.tryAdd(message); - if (!added) { - System.out.println(getClass().getName() + " !!!! ERROR could not add index =" ); - throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, "Could not fit in one batch message", - null); - } - } - ); - return connectionProcessor - .flatMap(connection -> connection.getManagementNode(entityName, entityType)) - .flatMapMany(managementNode -> managementNode.schedule(messageBatch.getMessages(), scheduledEnqueueTime, - messageBatch.getMaxSizeInBytes(), link.getLinkName(), transactionContext)); - })); - - - /* final CreateBatchOptions batchOptions = new CreateBatchOptions() - .setMaximumSizeInBytes(batchSize); - return messages.collect(new AmqpMessageCollector(batchOptions, 1, - link::getErrorContext, tracerProvider, messageSerializer, entityName, - link.getHostname()));*/ - })); - + return getSendLink().flatMapMany(link -> link.getLinkSize() + .flatMapMany(size -> { + int maxSize = size > 0 + ? size + : MAX_MESSAGE_LENGTH_BYTES; + final CreateBatchOptions batchOptions = new CreateBatchOptions(); + batchOptions.setMaximumSizeInBytes(maxSize); + + return createBatch(batchOptions).flatMapMany(messageBatch -> { + AtomicInteger index = new AtomicInteger(); + messages.forEach(message -> { + index.incrementAndGet(); + boolean added = messageBatch.tryAdd(message); + if (!added) { + final String error = String.format(Locale.US, + "Messages exceed max allowed size. Failed to add message at index '%s'.", + index.get()); + throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, error, + link.getErrorContext()); + } + }); + + return connectionProcessor + .flatMap(connection -> connection.getManagementNode(entityName, entityType)) + .flatMapMany(managementNode -> managementNode.schedule(messageBatch.getMessages(), + scheduledEnqueueTime, messageBatch.getMaxSizeInBytes(), link.getLinkName(), + transactionContext)); + }); + })); } - /*{ - if (!messageBatch.tryAdd(message)) { - - final String errorMessage = String.format(Locale.US, - "EventData does not fit into maximum number of batches. '%s'", maxNumberOfBatches); - - throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, errorMessage, - contextProvider.getErrorContext()); - - };*/ /** * Cancels the enqueuing of an already scheduled message, if it was not already enqueued. * diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java index ec9048938385..f7465c97ccfb 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java @@ -58,7 +58,6 @@ import java.util.Collections; import java.util.List; import java.util.UUID; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.IntStream; @@ -71,6 +70,8 @@ import static com.azure.messaging.servicebus.implementation.ServiceBusConstants.AZ_TRACING_SERVICE_NAME; import static java.nio.charset.StandardCharsets.UTF_8; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyList; import static org.mockito.ArgumentMatchers.anyString; @@ -248,10 +249,8 @@ void createBatchWhenSizeTooBig() { } @Test - void scheduleMessageSizeTooBig() throws InterruptedException { + void scheduleMessageSizeTooBig() { // Arrange - long sequenceNumberReturned = 10; - int maxLinkSize = 1024; int batchSize = maxLinkSize + 10; @@ -264,22 +263,14 @@ void scheduleMessageSizeTooBig() throws InterruptedException { when(connection.createSendLink(eq(ENTITY_NAME), eq(ENTITY_NAME), any(AmqpRetryOptions.class), isNull())) .thenReturn(Mono.just(link)); when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); - when(managementNode.schedule(anyList(), eq(instant), any(Integer.class), any(), isNull())) - .thenReturn(Flux.just(sequenceNumberReturned)); // Act & Assert - sender.scheduleMessages(messages, instant) - .subscribe(aLong -> { - System.out.println("!!!! Test Messages Scheduled : " + aLong); + StepVerifier.create(sender.scheduleMessages(messages, instant)) + .verifyErrorMatches(throwable -> { + assertTrue(throwable instanceof AmqpException); + assertSame(((AmqpException) throwable).getErrorCondition(), AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED); + return true; }); - TimeUnit.SECONDS.sleep(5); - - verify(managementNode).schedule(sbMessagesCaptor.capture(), eq(instant), eq(MAX_MESSAGE_LENGTH_BYTES), eq(LINK_NAME), isNull()); - List actualMessages = sbMessagesCaptor.getValue(); - Assertions.assertNotNull(actualMessages); - System.out.println(" !!!! Test actualMessages : " + actualMessages.size()); - //Assertions.assertEquals(1, actualMessages.size()); - //Assertions.assertEquals(message, actualMessages.get(0)); } /** diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java index abf9bb236cea..4f8fc0677b5d 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java @@ -346,7 +346,7 @@ void multipleReceiveUnnamedSession() { // Arrange final int expectedLinksCreated = 2; final ReceiverOptions receiverOptions = new ReceiverOptions(ReceiveMode.PEEK_LOCK, 1, null, - false, 1, Duration.ZERO); + true, 1, Duration.ZERO); sessionManager = new ServiceBusSessionManager(ENTITY_PATH, ENTITY_TYPE, connectionProcessor, tracerProvider, messageSerializer, receiverOptions); From a30cc197b877b2b09ba452d649b0aa15ef0bb32f Mon Sep 17 00:00:00 2001 From: Hemant Tanwar Date: Fri, 23 Oct 2020 11:37:56 -0700 Subject: [PATCH 3/6] added check for maxsize for schedule message api --- .../servicebus/ServiceBusSenderAsyncClient.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java index f4b2f7f2959c..b41b392b024d 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java @@ -339,7 +339,7 @@ public Mono scheduleMessage(ServiceBusMessage message, OffsetDateTime sche * Sends a scheduled messages to the Azure Service Bus entity this sender is connected to. A scheduled message is * enqueued and made available to receivers only at the scheduled enqueue time. * - * @param messages Message to be sent to the Service Bus Queue. + * @param messages Messages to be sent to the Service Bus Queue. * @param scheduledEnqueueTime OffsetDateTime at which the message should appear in the Service Bus queue or topic. * * @return The sequence number of the scheduled message which can be used to cancel the scheduling of the message. @@ -354,7 +354,7 @@ public Flux scheduleMessages(Iterable messages, OffsetD * Sends a scheduled messages to the Azure Service Bus entity this sender is connected to. A scheduled message is * enqueued and made available to receivers only at the scheduled enqueue time. * - * @param messages Message to be sent to the Service Bus Queue. + * @param messages Messages to be sent to the Service Bus Queue. * @param scheduledEnqueueTime Instant at which the message should appear in the Service Bus queue or topic. * @param transactionContext to be set on batch message before scheduling them on Service Bus. * @@ -387,10 +387,12 @@ public Flux scheduleMessages(Iterable messages, OffsetD boolean added = messageBatch.tryAdd(message); if (!added) { final String error = String.format(Locale.US, - "Messages exceed max allowed size. Failed to add message at index '%s'.", - index.get()); - throw new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, error, - link.getErrorContext()); + "Messages exceed max allowed size '%s' bytes for all the messages together." + + " Failed to add message at index '%s'.", + maxSize, index.get()); + //throw new ; + throw logger.logExceptionAsError(new AmqpException(false, + AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, error, link.getErrorContext())); } }); From c0942818ba672f487441c7ece954afb8352fe996 Mon Sep 17 00:00:00 2001 From: Hemant Tanwar Date: Fri, 23 Oct 2020 11:44:31 -0700 Subject: [PATCH 4/6] added check for maxsize for schedule message api --- .../ServiceBusSenderAsyncClient.java | 1 - .../ServiceBusSenderAsyncClientTest.java | 50 +++++++++---------- 2 files changed, 25 insertions(+), 26 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java index b41b392b024d..592869c3777c 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java @@ -390,7 +390,6 @@ public Flux scheduleMessages(Iterable messages, OffsetD "Messages exceed max allowed size '%s' bytes for all the messages together." + " Failed to add message at index '%s'.", maxSize, index.get()); - //throw new ; throw logger.logExceptionAsError(new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, error, link.getErrorContext())); } diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java index f7465c97ccfb..cf4e954163a7 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientTest.java @@ -248,31 +248,6 @@ void createBatchWhenSizeTooBig() { .verify(); } - @Test - void scheduleMessageSizeTooBig() { - // Arrange - int maxLinkSize = 1024; - int batchSize = maxLinkSize + 10; - - OffsetDateTime instant = mock(OffsetDateTime.class); - final List messages = TestUtils.getServiceBusMessages(batchSize, UUID.randomUUID().toString()); - - final AmqpSendLink link = mock(AmqpSendLink.class); - when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); - - when(connection.createSendLink(eq(ENTITY_NAME), eq(ENTITY_NAME), any(AmqpRetryOptions.class), isNull())) - .thenReturn(Mono.just(link)); - when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); - - // Act & Assert - StepVerifier.create(sender.scheduleMessages(messages, instant)) - .verifyErrorMatches(throwable -> { - assertTrue(throwable instanceof AmqpException); - assertSame(((AmqpException) throwable).getErrorCondition(), AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED); - return true; - }); - } - /** * Verifies that the producer can create a batch with a given {@link CreateBatchOptions#getMaximumSizeInBytes()}. */ @@ -315,6 +290,31 @@ void createsMessageBatchWithSize() { .verifyComplete(); } + @Test + void scheduleMessageSizeTooBig() { + // Arrange + int maxLinkSize = 1024; + int batchSize = maxLinkSize + 10; + + OffsetDateTime instant = mock(OffsetDateTime.class); + final List messages = TestUtils.getServiceBusMessages(batchSize, UUID.randomUUID().toString()); + + final AmqpSendLink link = mock(AmqpSendLink.class); + when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); + + when(connection.createSendLink(eq(ENTITY_NAME), eq(ENTITY_NAME), any(AmqpRetryOptions.class), isNull())) + .thenReturn(Mono.just(link)); + when(link.getLinkSize()).thenReturn(Mono.just(maxLinkSize)); + + // Act & Assert + StepVerifier.create(sender.scheduleMessages(messages, instant)) + .verifyErrorMatches(throwable -> { + assertTrue(throwable instanceof AmqpException); + assertSame(((AmqpException) throwable).getErrorCondition(), AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED); + return true; + }); + } + /** * Verifies that sending multiple message will result in calling sender.send(MessageBatch, transaction). */ From 8dcea5b7c715c20eacf882db862e5ab55083a034 Mon Sep 17 00:00:00 2001 From: Hemant Tanwar Date: Fri, 23 Oct 2020 12:31:27 -0700 Subject: [PATCH 5/6] added check for maxsize for schedule message api --- .../messaging/servicebus/ServiceBusSessionManagerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java index 4f8fc0677b5d..abf9bb236cea 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionManagerTest.java @@ -346,7 +346,7 @@ void multipleReceiveUnnamedSession() { // Arrange final int expectedLinksCreated = 2; final ReceiverOptions receiverOptions = new ReceiverOptions(ReceiveMode.PEEK_LOCK, 1, null, - true, 1, Duration.ZERO); + false, 1, Duration.ZERO); sessionManager = new ServiceBusSessionManager(ENTITY_PATH, ENTITY_TYPE, connectionProcessor, tracerProvider, messageSerializer, receiverOptions); From 7b3e2f35084c39f26d5827b8e9add7e085624282 Mon Sep 17 00:00:00 2001 From: Hemant Tanwar Date: Wed, 28 Oct 2020 15:49:01 -0700 Subject: [PATCH 6/6] Review comments --- .../messaging/servicebus/ServiceBusSenderAsyncClient.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java index 66a29346b655..3332ced69344 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClient.java @@ -375,15 +375,14 @@ public Flux scheduleMessages(Iterable messages, OffsetD .flatMapMany(messageBatch -> { int index = 0; for (ServiceBusMessage message : messages) { - ++index; - boolean added = messageBatch.tryAddMessage(message); - if (!added) { + if (!messageBatch.tryAddMessage(message)) { final String error = String.format(Locale.US, "Messages exceed max allowed size for all the messages together. " + "Failed to add message at index '%s'.", index); throw logger.logExceptionAsError(new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED, error, link.getErrorContext())); } + ++index; } return connectionProcessor