From 75f14a5d23a245af90a189c4475ab36b71df39a8 Mon Sep 17 00:00:00 2001 From: Ramya Achutha Rao Date: Fri, 30 Oct 2020 11:30:22 -0700 Subject: [PATCH 1/3] [Service Bus] Remove send via option from sender builder --- .../servicebus/ServiceBusClientBuilder.java | 44 +-- .../ServiceBusClientBuilderTest.java | 60 ---- ...ceBusSenderAsyncClientIntegrationTest.java | 290 +++++++++--------- 3 files changed, 147 insertions(+), 247 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java index 98b47ad619e0..4b68c5608bd4 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java @@ -494,8 +494,6 @@ private static String getEntityPath(ClientLogger logger, MessagingEntityType ent public final class ServiceBusSenderClientBuilder { private String queueName; private String topicName; - private String viaQueueName; - private String viaTopicName; private ServiceBusSenderClientBuilder() { } @@ -512,34 +510,6 @@ public ServiceBusSenderClientBuilder queueName(String queueName) { return this; } - /** - * Sets the name of the initial destination Service Bus queue to publish messages to. - * - * @param viaQueueName The initial destination of the message. - * - * @return The modified {@link ServiceBusSenderClientBuilder} object. - * @see Send - * Via - */ - public ServiceBusSenderClientBuilder viaQueueName(String viaQueueName) { - this.viaQueueName = viaQueueName; - return this; - } - - /** - * Sets the name of the initial destination Service Bus topic to publish messages to. - * - * @param viaTopicName The initial destination of the message. - * - * @return The modified {@link ServiceBusSenderClientBuilder} object. - * @see Send - * Via - */ - public ServiceBusSenderClientBuilder viaTopicName(String viaTopicName) { - this.viaTopicName = viaTopicName; - return this; - } - /** * Sets the name of the Service Bus topic to publish messages to. * @@ -560,8 +530,7 @@ public ServiceBusSenderClientBuilder topicName(String topicName) { * @throws IllegalStateException if {@link #queueName(String) queueName} or {@link #topicName(String) * topicName} are not set or, both of these fields are set. It is also thrown if the Service Bus {@link * #connectionString(String) connectionString} contains an {@code EntityPath} that does not match one set in - * {@link #queueName(String) queueName} or {@link #topicName(String) topicName}. Or the {@link - * #viaQueueName(String) viaQueueName} is specified along with {@link #topicName(String) topicName}. + * {@link #queueName(String) queueName} or {@link #topicName(String) topicName}. * @throws IllegalArgumentException if the entity type is not a queue or a topic. */ public ServiceBusSenderAsyncClient buildAsyncClient() { @@ -569,16 +538,7 @@ public ServiceBusSenderAsyncClient buildAsyncClient() { final MessagingEntityType entityType = validateEntityPaths(logger, connectionStringEntityName, topicName, queueName); - if (!CoreUtils.isNullOrEmpty(viaQueueName) && entityType == MessagingEntityType.SUBSCRIPTION) { - throw logger.logExceptionAsError(new IllegalStateException(String.format( - "(%s), Via queue feature work only with a queue.", viaQueueName))); - } else if (!CoreUtils.isNullOrEmpty(viaTopicName) && entityType == MessagingEntityType.QUEUE) { - throw logger.logExceptionAsError(new IllegalStateException(String.format( - "(%s), Via topic feature work only with a topic.", viaTopicName))); - } - final String entityName; - final String viaEntityName = !CoreUtils.isNullOrEmpty(viaQueueName) ? viaQueueName : viaTopicName; switch (entityType) { case QUEUE: entityName = queueName; @@ -595,7 +555,7 @@ public ServiceBusSenderAsyncClient buildAsyncClient() { } return new ServiceBusSenderAsyncClient(entityName, entityType, connectionProcessor, retryOptions, - tracerProvider, messageSerializer, ServiceBusClientBuilder.this::onClientClose, viaEntityName); + tracerProvider, messageSerializer, ServiceBusClientBuilder.this::onClientClose, null); } /** diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusClientBuilderTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusClientBuilderTest.java index 71d08def3139..fae12997fb9c 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusClientBuilderTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusClientBuilderTest.java @@ -52,66 +52,6 @@ class ServiceBusClientBuilderTest { ENDPOINT, SHARED_ACCESS_KEY_NAME, SHARED_ACCESS_KEY, QUEUE_NAME); private static final Proxy PROXY_ADDRESS = new Proxy(Proxy.Type.HTTP, new InetSocketAddress(PROXY_HOST, Integer.parseInt(PROXY_PORT))); - @Test - void viaQueueNameWithTopicNotAllowed() { - // Arrange - ServiceBusSenderClientBuilder builder = new ServiceBusClientBuilder() - .connectionString(NAMESPACE_CONNECTION_STRING) - .sender() - .topicName(TOPIC_NAME) - .viaQueueName(VIA_QUEUE_NAME); - - // Act & Assert - assertThrows(IllegalStateException.class, () -> builder.buildAsyncClient()); - } - - @Test - void viaTopicNameWithQueueNotAllowed() { - // Arrange - ServiceBusSenderClientBuilder builder = new ServiceBusClientBuilder() - .connectionString(NAMESPACE_CONNECTION_STRING) - .sender() - .queueName(QUEUE_NAME) - .viaTopicName(VIA_TOPIC_NAME); - - // Act & Assert - assertThrows(IllegalStateException.class, () -> builder.buildAsyncClient()); - } - - @Test - void queueClientWithViaQueueName() { - // Arrange - final ServiceBusSenderClientBuilder builder = new ServiceBusClientBuilder() - .connectionString(NAMESPACE_CONNECTION_STRING) - .sender() - .queueName(QUEUE_NAME) - .viaQueueName(VIA_QUEUE_NAME); - - // Act - final ServiceBusSenderAsyncClient client = builder.buildAsyncClient(); - - // Assert - assertNotNull(client); - assertEquals(client.getEntityPath(), QUEUE_NAME); - } - - @Test - void topicClientWithViaTopicName() { - // Arrange - final ServiceBusSenderClientBuilder builder = new ServiceBusClientBuilder() - .connectionString(NAMESPACE_CONNECTION_STRING) - .sender() - .topicName(TOPIC_NAME) - .viaTopicName(VIA_TOPIC_NAME); - - // Act - final ServiceBusSenderAsyncClient client = builder.buildAsyncClient(); - - // Assert - assertNotNull(client); - assertEquals(client.getEntityPath(), TOPIC_NAME); - } - @Test void deadLetterqueueClient() { // Arrange diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java index 3adf24e6822c..a5b9214185e1 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java @@ -147,156 +147,156 @@ void nonSessionMessageBatch(MessagingEntityType entityType) { * Verifies that we can send message to final destination using via-queue. */ @Test - void viaQueueMessageSendTest() { - // Arrange - final boolean useCredentials = false; - final Duration shortTimeout = Duration.ofSeconds(15); - final int viaIntermediateEntity = TestUtils.USE_CASE_SEND_VIA_QUEUE_1; - final int destinationEntity = TestUtils.USE_CASE_SEND_VIA_QUEUE_2; - final boolean shareConnection = true; - final MessagingEntityType entityType = MessagingEntityType.QUEUE; - final boolean isSessionEnabled = false; - final String messageId = UUID.randomUUID().toString(); - final int total = 1; - final int totalToDestination = 2; - final List messages = TestUtils.getServiceBusMessages(total, messageId, CONTENTS_BYTES); - final String viaQueueName = getQueueName(viaIntermediateEntity); - - setSenderAndReceiver(entityType, viaIntermediateEntity, useCredentials); - - final ServiceBusSenderAsyncClient destination1ViaSender = getSenderBuilder(useCredentials, entityType, - destinationEntity, false, shareConnection) - .viaQueueName(viaQueueName) - .buildAsyncClient(); - final ServiceBusReceiverAsyncClient destination1Receiver = getReceiverBuilder(useCredentials, entityType, - destinationEntity, shareConnection) - .receiveMode(ReceiveMode.RECEIVE_AND_DELETE) - .buildAsyncClient(); - - final AtomicReference transaction = new AtomicReference<>(); - - // Act - try { - StepVerifier.create(destination1ViaSender.createTransaction()) - .assertNext(transactionContext -> { - transaction.set(transactionContext); - assertNotNull(transaction); - }) - .verifyComplete(); - assertNotNull(transaction.get()); - - StepVerifier.create(sender.sendMessages(messages, transaction.get())) - .verifyComplete(); - StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) - .verifyComplete(); - StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) - .verifyComplete(); - - StepVerifier.create(destination1ViaSender.commitTransaction(transaction.get()) - .delaySubscription(Duration.ofSeconds(1))) - .verifyComplete(); - - // Assert - // Verify message is received by final destination Entity - StepVerifier.create(destination1Receiver.receiveMessages().take(totalToDestination).timeout(shortTimeout)) - .assertNext(receivedMessage -> { - assertMessageEquals(receivedMessage, messageId, isSessionEnabled); - messagesPending.decrementAndGet(); - }) - .assertNext(receivedMessage -> { - assertMessageEquals(receivedMessage, messageId, isSessionEnabled); - messagesPending.decrementAndGet(); - }) - .verifyComplete(); - - // Verify, intermediate-via queue has it delivered to intermediate Entity. - StepVerifier.create(receiver.receiveMessages().take(total).timeout(shortTimeout)) - .assertNext(receivedMessage -> { - assertMessageEquals(receivedMessage, messageId, isSessionEnabled); - messagesPending.decrementAndGet(); - }) - .verifyComplete(); - } finally { - destination1Receiver.close(); - destination1ViaSender.close(); - } - } + // void viaQueueMessageSendTest() { + // // Arrange + // final boolean useCredentials = false; + // final Duration shortTimeout = Duration.ofSeconds(15); + // final int viaIntermediateEntity = TestUtils.USE_CASE_SEND_VIA_QUEUE_1; + // final int destinationEntity = TestUtils.USE_CASE_SEND_VIA_QUEUE_2; + // final boolean shareConnection = true; + // final MessagingEntityType entityType = MessagingEntityType.QUEUE; + // final boolean isSessionEnabled = false; + // final String messageId = UUID.randomUUID().toString(); + // final int total = 1; + // final int totalToDestination = 2; + // final List messages = TestUtils.getServiceBusMessages(total, messageId, CONTENTS_BYTES); + // final String viaQueueName = getQueueName(viaIntermediateEntity); + + // setSenderAndReceiver(entityType, viaIntermediateEntity, useCredentials); + + // final ServiceBusSenderAsyncClient destination1ViaSender = getSenderBuilder(useCredentials, entityType, + // destinationEntity, false, shareConnection) + // .viaQueueName(viaQueueName) + // .buildAsyncClient(); + // final ServiceBusReceiverAsyncClient destination1Receiver = getReceiverBuilder(useCredentials, entityType, + // destinationEntity, shareConnection) + // .receiveMode(ReceiveMode.RECEIVE_AND_DELETE) + // .buildAsyncClient(); + + // final AtomicReference transaction = new AtomicReference<>(); + + // // Act + // try { + // StepVerifier.create(destination1ViaSender.createTransaction()) + // .assertNext(transactionContext -> { + // transaction.set(transactionContext); + // assertNotNull(transaction); + // }) + // .verifyComplete(); + // assertNotNull(transaction.get()); + + // StepVerifier.create(sender.sendMessages(messages, transaction.get())) + // .verifyComplete(); + // StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) + // .verifyComplete(); + // StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) + // .verifyComplete(); + + // StepVerifier.create(destination1ViaSender.commitTransaction(transaction.get()) + // .delaySubscription(Duration.ofSeconds(1))) + // .verifyComplete(); + + // // Assert + // // Verify message is received by final destination Entity + // StepVerifier.create(destination1Receiver.receiveMessages().take(totalToDestination).timeout(shortTimeout)) + // .assertNext(receivedMessage -> { + // assertMessageEquals(receivedMessage, messageId, isSessionEnabled); + // messagesPending.decrementAndGet(); + // }) + // .assertNext(receivedMessage -> { + // assertMessageEquals(receivedMessage, messageId, isSessionEnabled); + // messagesPending.decrementAndGet(); + // }) + // .verifyComplete(); + + // // Verify, intermediate-via queue has it delivered to intermediate Entity. + // StepVerifier.create(receiver.receiveMessages().take(total).timeout(shortTimeout)) + // .assertNext(receivedMessage -> { + // assertMessageEquals(receivedMessage, messageId, isSessionEnabled); + // messagesPending.decrementAndGet(); + // }) + // .verifyComplete(); + // } finally { + // destination1Receiver.close(); + // destination1ViaSender.close(); + // } + // } /** * Verifies that we can send message to final destination using via-topic. */ @Test - void viaTopicMessageSendTest() { - // Arrange - final boolean useCredentials = false; - final Duration shortTimeout = Duration.ofSeconds(15); - final int viaIntermediateEntity = TestUtils.USE_CASE_SEND_VIA_TOPIC_1; - final int destinationEntity = TestUtils.USE_CASE_SEND_VIA_TOPIC_2; - final boolean shareConnection = true; - final MessagingEntityType entityType = MessagingEntityType.SUBSCRIPTION; - final boolean isSessionEnabled = false; - final String messageId = UUID.randomUUID().toString(); - final int total = 1; - final int totalToDestination = 2; - final List messages = TestUtils.getServiceBusMessages(total, messageId, CONTENTS_BYTES); - final String viaTopicName = getTopicName(viaIntermediateEntity); - - setSenderAndReceiver(entityType, viaIntermediateEntity, useCredentials); - final ServiceBusReceiverAsyncClient intermediateReceiver = receiver; - final ServiceBusSenderAsyncClient intermediateSender = sender; - - final ServiceBusSenderAsyncClient destination1ViaSender = getSenderBuilder(useCredentials, entityType, - destinationEntity, false, shareConnection) - .viaTopicName(viaTopicName) - .buildAsyncClient(); - - final ServiceBusReceiverAsyncClient destination1Receiver = getReceiverBuilder(useCredentials, entityType, - destinationEntity, shareConnection) - .receiveMode(ReceiveMode.RECEIVE_AND_DELETE) - .buildAsyncClient(); - - final AtomicReference transaction = new AtomicReference<>(); - - // Act - StepVerifier.create(destination1ViaSender.createTransaction()) - .assertNext(transactionContext -> { - transaction.set(transactionContext); - assertNotNull(transaction); - }) - .verifyComplete(); - assertNotNull(transaction.get()); - - StepVerifier.create(intermediateSender.sendMessages(messages, transaction.get())) - .verifyComplete(); - StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) - .verifyComplete(); - StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) - .verifyComplete(); - - StepVerifier.create(destination1ViaSender.commitTransaction(transaction.get()).delaySubscription(Duration.ofSeconds(1))) - .verifyComplete(); - - // Assert - // Verify message is received by final destination Entity - StepVerifier.create(destination1Receiver.receiveMessages().take(totalToDestination).timeout(shortTimeout)) - .assertNext(receivedMessage -> { - assertMessageEquals(receivedMessage, messageId, isSessionEnabled); - messagesPending.decrementAndGet(); - }) - .assertNext(receivedMessage -> { - assertMessageEquals(receivedMessage, messageId, isSessionEnabled); - messagesPending.decrementAndGet(); - }) - .verifyComplete(); - - // Verify, intermediate-via topic has it delivered to intermediate Entity. - StepVerifier.create(intermediateReceiver.receiveMessages().take(total).timeout(shortTimeout)) - .assertNext(receivedMessage -> { - assertMessageEquals(receivedMessage, messageId, isSessionEnabled); - messagesPending.decrementAndGet(); - }) - .verifyComplete(); - } + // void viaTopicMessageSendTest() { + // // Arrange + // final boolean useCredentials = false; + // final Duration shortTimeout = Duration.ofSeconds(15); + // final int viaIntermediateEntity = TestUtils.USE_CASE_SEND_VIA_TOPIC_1; + // final int destinationEntity = TestUtils.USE_CASE_SEND_VIA_TOPIC_2; + // final boolean shareConnection = true; + // final MessagingEntityType entityType = MessagingEntityType.SUBSCRIPTION; + // final boolean isSessionEnabled = false; + // final String messageId = UUID.randomUUID().toString(); + // final int total = 1; + // final int totalToDestination = 2; + // final List messages = TestUtils.getServiceBusMessages(total, messageId, CONTENTS_BYTES); + // final String viaTopicName = getTopicName(viaIntermediateEntity); + + // setSenderAndReceiver(entityType, viaIntermediateEntity, useCredentials); + // final ServiceBusReceiverAsyncClient intermediateReceiver = receiver; + // final ServiceBusSenderAsyncClient intermediateSender = sender; + + // final ServiceBusSenderAsyncClient destination1ViaSender = getSenderBuilder(useCredentials, entityType, + // destinationEntity, false, shareConnection) + // .viaTopicName(viaTopicName) + // .buildAsyncClient(); + + // final ServiceBusReceiverAsyncClient destination1Receiver = getReceiverBuilder(useCredentials, entityType, + // destinationEntity, shareConnection) + // .receiveMode(ReceiveMode.RECEIVE_AND_DELETE) + // .buildAsyncClient(); + + // final AtomicReference transaction = new AtomicReference<>(); + + // // Act + // StepVerifier.create(destination1ViaSender.createTransaction()) + // .assertNext(transactionContext -> { + // transaction.set(transactionContext); + // assertNotNull(transaction); + // }) + // .verifyComplete(); + // assertNotNull(transaction.get()); + + // StepVerifier.create(intermediateSender.sendMessages(messages, transaction.get())) + // .verifyComplete(); + // StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) + // .verifyComplete(); + // StepVerifier.create(destination1ViaSender.sendMessages(messages, transaction.get())) + // .verifyComplete(); + + // StepVerifier.create(destination1ViaSender.commitTransaction(transaction.get()).delaySubscription(Duration.ofSeconds(1))) + // .verifyComplete(); + + // // Assert + // // Verify message is received by final destination Entity + // StepVerifier.create(destination1Receiver.receiveMessages().take(totalToDestination).timeout(shortTimeout)) + // .assertNext(receivedMessage -> { + // assertMessageEquals(receivedMessage, messageId, isSessionEnabled); + // messagesPending.decrementAndGet(); + // }) + // .assertNext(receivedMessage -> { + // assertMessageEquals(receivedMessage, messageId, isSessionEnabled); + // messagesPending.decrementAndGet(); + // }) + // .verifyComplete(); + + // // Verify, intermediate-via topic has it delivered to intermediate Entity. + // StepVerifier.create(intermediateReceiver.receiveMessages().take(total).timeout(shortTimeout)) + // .assertNext(receivedMessage -> { + // assertMessageEquals(receivedMessage, messageId, isSessionEnabled); + // messagesPending.decrementAndGet(); + // }) + // .verifyComplete(); + // } /** * Verifies that we can do following From c08f9dec0cc669cf671ff1abd22c24b51dcb9ea5 Mon Sep 17 00:00:00 2001 From: Ramya Achutha Rao Date: Fri, 30 Oct 2020 11:49:20 -0700 Subject: [PATCH 2/3] Comment the test annotation --- .../ServiceBusSenderAsyncClientIntegrationTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java index a5b9214185e1..df570497cd84 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java @@ -146,7 +146,7 @@ void nonSessionMessageBatch(MessagingEntityType entityType) { /** * Verifies that we can send message to final destination using via-queue. */ - @Test + // @Test // void viaQueueMessageSendTest() { // // Arrange // final boolean useCredentials = false; @@ -225,7 +225,7 @@ void nonSessionMessageBatch(MessagingEntityType entityType) { /** * Verifies that we can send message to final destination using via-topic. */ - @Test + // @Test // void viaTopicMessageSendTest() { // // Arrange // final boolean useCredentials = false; From fd8cb5e0350bf193646b3eee18d44a8ddf9933f5 Mon Sep 17 00:00:00 2001 From: Ramya Achutha Rao Date: Fri, 30 Oct 2020 12:01:53 -0700 Subject: [PATCH 3/3] Remove unused import --- .../servicebus/ServiceBusSenderAsyncClientIntegrationTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java index df570497cd84..9018c802ede0 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSenderAsyncClientIntegrationTest.java @@ -9,7 +9,6 @@ import com.azure.messaging.servicebus.models.ReceiveMode; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Tag; -import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource;