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..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;
@@ -146,157 +145,157 @@ 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();
- }
- }
+ // @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();
+ // }
+ // }
/**
* 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();
- }
+ // @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();
+ // }
/**
* Verifies that we can do following