Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,12 @@ protected void afterTest() {
@Test
public void backCompatWithJavaSDKOlderThan0110() {
// Arrange
final Duration timeout = Duration.ofSeconds(30);
final String messageTrackingValue = UUID.randomUUID().toString();
final PartitionProperties properties = consumer.getPartitionProperties(PARTITION_ID).block(timeout);

Assertions.assertNotNull(properties);
final EventPosition position = EventPosition.fromSequenceNumber(properties.getLastEnqueuedSequenceNumber(), true);

// until version 0.10.0 - we used to have Properties as HashMap<String,String>
// This specific combination is intended to test the back compat - with the new Properties type as HashMap<String, Object>
Expand All @@ -96,12 +101,13 @@ public void backCompatWithJavaSDKOlderThan0110() {
final EventData eventData = serializer.deserialize(message, EventData.class);

// Act & Assert
StepVerifier.create(consumer.receiveFromPartition(PARTITION_ID, EventPosition.latest())
producer.send(eventData, sendOptions).block(TIMEOUT);

StepVerifier.create(consumer.receiveFromPartition(PARTITION_ID, position)
.filter(received -> isMatchingEvent(received, messageTrackingValue)).take(1))
.then(() -> producer.send(eventData, sendOptions).block(TIMEOUT))
.assertNext(event -> validateAmqpProperties(applicationProperties, event.getData()))
.expectComplete()
.verify(Duration.ofSeconds(45));
.verify(timeout);
}

private void validateAmqpProperties(Map<String, Object> expected, EventData event) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,7 @@ void getPropertiesWithCredentials() {
StepVerifier.create(client.getProperties())
.assertNext(properties -> {
Assertions.assertEquals(getEventHubName(), properties.getName());
Assertions.assertEquals(2, properties.getPartitionIds().stream().count());
Assertions.assertEquals(3, properties.getPartitionIds().stream().count());
})
.expectComplete()
.verify(TIMEOUT);
Expand All @@ -189,7 +189,7 @@ void getMultipleProperties() {
StepVerifier.create(theClient.getProperties())
.assertNext(properties -> {
Assertions.assertEquals(getEventHubName(), properties.getName());
Assertions.assertEquals(2, properties.getPartitionIds().stream().count());
Assertions.assertEquals(3, properties.getPartitionIds().stream().count());
})
.verifyComplete();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
* Tests the metadata operations such as fetching partition properties and event hub properties.
*/
public class EventHubClientMetadataIntegrationTest extends IntegrationTestBase {
private final String[] expectedPartitionIds = new String[]{"0", "1"};
private final String[] expectedPartitionIds = new String[]{"0", "1", "2"};
private EventHubAsyncClient client;
private String eventHubName;

Expand Down Expand Up @@ -92,6 +92,7 @@ public void getPartitionPropertiesMultipleCalls() {

// Assert
StepVerifier.create(partitionProperties)
.assertNext(properties -> Assertions.assertEquals(eventHubName, properties.getEventHubName()))
.assertNext(properties -> Assertions.assertEquals(eventHubName, properties.getEventHubName()))
.assertNext(properties -> Assertions.assertEquals(eventHubName, properties.getEventHubName()))
.verifyComplete();
Expand All @@ -111,16 +112,20 @@ public void getPartitionPropertiesInvalidToken() {
.buildAsyncClient();

// Act & Assert
StepVerifier.create(invalidClient.getProperties())
.expectErrorSatisfies(error -> {
Assertions.assertTrue(error instanceof AmqpException);

AmqpException exception = (AmqpException) error;
Assertions.assertEquals(AmqpErrorCondition.UNAUTHORIZED_ACCESS, exception.getErrorCondition());
Assertions.assertFalse(exception.isTransient());
Assertions.assertFalse(CoreUtils.isNullOrEmpty(exception.getMessage()));
})
.verify();
try {
StepVerifier.create(invalidClient.getProperties())
.expectErrorSatisfies(error -> {
Assertions.assertTrue(error instanceof AmqpException);

AmqpException exception = (AmqpException) error;
Assertions.assertEquals(AmqpErrorCondition.UNAUTHORIZED_ACCESS, exception.getErrorCondition());
Assertions.assertFalse(exception.isTransient());
Assertions.assertFalse(CoreUtils.isNullOrEmpty(exception.getMessage()));
})
.verify();
} finally {
invalidClient.close();
}
}

/**
Expand All @@ -137,15 +142,19 @@ public void getPartitionPropertiesNonExistentHub() {
.buildAsyncClient();

// Act & Assert
StepVerifier.create(invalidClient.getPartitionIds())
.expectErrorSatisfies(error -> {
Assertions.assertTrue(error instanceof AmqpException);

AmqpException exception = (AmqpException) error;
Assertions.assertEquals(AmqpErrorCondition.NOT_FOUND, exception.getErrorCondition());
Assertions.assertFalse(exception.isTransient());
Assertions.assertFalse(CoreUtils.isNullOrEmpty(exception.getMessage()));
})
.verify();
try {
StepVerifier.create(invalidClient.getPartitionIds())
.expectErrorSatisfies(error -> {
Assertions.assertTrue(error instanceof AmqpException);

AmqpException exception = (AmqpException) error;
Assertions.assertEquals(AmqpErrorCondition.NOT_FOUND, exception.getErrorCondition());
Assertions.assertFalse(exception.isTransient());
Assertions.assertFalse(CoreUtils.isNullOrEmpty(exception.getMessage()));
})
.verify();
} finally {
invalidClient.close();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,10 @@
import com.azure.messaging.eventhubs.models.PartitionEvent;
import com.azure.messaging.eventhubs.models.ReceiveOptions;
import com.azure.messaging.eventhubs.models.SendOptions;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import reactor.core.Disposable;
import reactor.core.Disposables;
Expand Down Expand Up @@ -46,7 +48,7 @@
public class EventHubConsumerAsyncClientIntegrationTest extends IntegrationTestBase {
private static final String PARTITION_ID_HEADER = "SENT_PARTITION_ID";

private final String[] expectedPartitionIds = new String[]{"0", "1"};
private final String[] expectedPartitionIds = new String[]{"0", "1", "2"};

private static final String MESSAGE_TRACKING_ID = UUID.randomUUID().toString();

Expand All @@ -56,6 +58,16 @@ public EventHubConsumerAsyncClientIntegrationTest() {
super(new ClientLogger(EventHubConsumerAsyncClientIntegrationTest.class));
}

@BeforeAll
static void beforeAll() {
StepVerifier.setDefaultTimeout(Duration.ofSeconds(30));
}

@AfterAll
static void afterAll() {
StepVerifier.resetDefaultTimeout();
}

@Override
protected void beforeTest() {
client = createBuilder()
Expand Down Expand Up @@ -138,21 +150,25 @@ public void parallelCreationOfReceivers() {
@Test
public void lastEnqueuedInformationIsNotUpdated() {
// Arrange
final String secondPartitionId = "1";
final EventPosition position = EventPosition.fromEnqueuedTime(Instant.now());
final String firstPartition = "0";
final PartitionProperties properties = client.getPartitionProperties(firstPartition).block();
Assertions.assertNotNull(properties);

final EventPosition position = EventPosition.fromSequenceNumber(properties.getLastEnqueuedSequenceNumber());
final EventHubConsumerAsyncClient consumer = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, 1);
final ReceiveOptions options = new ReceiveOptions().setTrackLastEnqueuedEventProperties(false);

final AtomicBoolean isActive = new AtomicBoolean(true);
final int expectedNumber = 5;
final EventHubProducerAsyncClient producer = client.createProducer();
final Disposable producerEvents = getEvents(isActive).flatMap(event -> producer.send(event)).subscribe(
sent -> logger.info("Event sent."),
error -> logger.error("Error sending event", error));
final SendOptions sendOptions = new SendOptions().setPartitionId(firstPartition);
final Disposable producerEvents = getEvents(isActive)
.flatMap(event -> producer.send(event, sendOptions))
.subscribe(sent -> logger.info("Event sent."), error -> logger.error("Error sending event", error));

// Act & Assert
try {
StepVerifier.create(consumer.receiveFromPartition(secondPartitionId, position, options)
StepVerifier.create(consumer.receiveFromPartition(firstPartition, position, options)
.take(expectedNumber))
.assertNext(event -> {
Assertions.assertNull(event.getLastEnqueuedEventProperties(), "'lastEnqueuedEventProperties' "
Expand Down Expand Up @@ -243,7 +259,8 @@ private static void verifyLastRetrieved(AtomicReference<LastEnqueuedEventPropert
public void sameOwnerLevelClosesFirstConsumer() throws InterruptedException {
// Arrange
final Semaphore semaphore = new Semaphore(1);
final String secondPartitionId = "1";

final String lastPartition = "2";
final EventPosition position = EventPosition.fromEnqueuedTime(Instant.now());
final ReceiveOptions options = new ReceiveOptions()
.setOwnerLevel(1L);
Expand All @@ -261,7 +278,7 @@ public void sameOwnerLevelClosesFirstConsumer() throws InterruptedException {
logger.info("STARTED CONSUMING FROM PARTITION 1");
semaphore.acquire();

subscriptions.add(consumer.receiveFromPartition(secondPartitionId, position)
subscriptions.add(consumer.receiveFromPartition(lastPartition, position)
.filter(event -> TestUtils.isMatchingEvent(event, MESSAGE_TRACKING_ID))
.subscribe(
event -> logger.info("C1:\tReceived event sequence: {}", event.getData().getSequenceNumber()),
Expand All @@ -277,7 +294,7 @@ public void sameOwnerLevelClosesFirstConsumer() throws InterruptedException {

logger.info("STARTED CONSUMING FROM PARTITION 1 with C3");
final EventHubConsumerAsyncClient consumer2 = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, 1);
subscriptions.add(consumer2.receiveFromPartition(secondPartitionId, position, options)
subscriptions.add(consumer2.receiveFromPartition(lastPartition, position, options)
.filter(event -> TestUtils.isMatchingEvent(event, MESSAGE_TRACKING_ID))
.subscribe(
event -> logger.info("C3:\tReceived event sequence: {}", event.getData().getSequenceNumber()),
Expand Down Expand Up @@ -310,7 +327,7 @@ public void getEventHubProperties() {
.assertNext(properties -> {
Assertions.assertNotNull(properties);
Assertions.assertEquals(consumer.getEventHubName(), properties.getName());
Assertions.assertEquals(2, properties.getPartitionIds().stream().count());
Assertions.assertEquals(3, properties.getPartitionIds().stream().count());
}).verifyComplete();
} finally {
dispose(consumer);
Expand Down Expand Up @@ -463,8 +480,8 @@ public void receivesMultiplePartitions() {
@Test
public void multipleReceiversSamePartition() throws InterruptedException {
// Arrange
final EventHubConsumerAsyncClient consumer = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, 50);
final EventHubConsumerAsyncClient consumer2 = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, 50);
final EventHubConsumerAsyncClient consumer = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, 1);
final EventHubConsumerAsyncClient consumer2 = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, 1);
final String partitionId = "1";
final PartitionProperties properties = consumer.getPartitionProperties(partitionId).block(TIMEOUT);
Assertions.assertNotNull(properties, "Should have been able to get partition properties.");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ public class EventHubConsumerClientIntegrationTest extends IntegrationTestBase {

private static final int NUMBER_OF_EVENTS = 10;
private static final AtomicBoolean HAS_PUSHED_EVENTS = new AtomicBoolean();
private final String[] expectedPartitionIds = new String[]{"0", "1"};
private final String[] expectedPartitionIds = new String[]{"0", "1", "2"};

private static volatile IntegrationTestEventData testData = null;

Expand Down Expand Up @@ -63,6 +63,8 @@ protected void beforeTest() {
testData = setupEventTestData(producer, NUMBER_OF_EVENTS, options);
}

Assertions.assertNotNull(testData, "'testData' should have been populated.");

startingPosition = EventPosition.fromEnqueuedTime(testData.getEnqueuedTime());
consumer = client.createConsumer(DEFAULT_CONSUMER_GROUP_NAME, DEFAULT_PREFETCH_COUNT);
}
Expand Down Expand Up @@ -242,7 +244,7 @@ public void getEventHubProperties() {
final EventHubProperties properties = consumer.getEventHubProperties();
Assertions.assertNotNull(properties);
Assertions.assertEquals(consumer.getEventHubName(), properties.getName());
Assertions.assertEquals(2, properties.getPartitionIds().stream().count());
Assertions.assertEquals(3, properties.getPartitionIds().stream().count());
} finally {
dispose(consumer);
}
Expand All @@ -262,7 +264,7 @@ public void getPartitionIds() {
final IterableStream<String> partitionIds = consumer.getPartitionIds();
final List<String> collect = partitionIds.stream().collect(Collectors.toList());

Assertions.assertEquals(2, collect.size());
Assertions.assertEquals(3, collect.size());
} finally {
dispose(consumer);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,7 @@ void sendWithCredentials() {
StepVerifier.create(client.getEventHubProperties())
.assertNext(properties -> {
Assertions.assertEquals(getEventHubName(), properties.getName());
Assertions.assertEquals(2, properties.getPartitionIds().stream().count());
Assertions.assertEquals(3, properties.getPartitionIds().stream().count());
})
.expectComplete()
.verify(TIMEOUT);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import com.azure.messaging.eventhubs.models.EventPosition;
import com.azure.messaging.eventhubs.models.SendOptions;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import reactor.core.publisher.Flux;
Expand All @@ -20,7 +21,7 @@
import java.net.ProxySelector;
import java.net.SocketAddress;
import java.net.URI;
import java.time.Instant;
import java.time.Duration;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
Expand All @@ -40,6 +41,8 @@ public ProxySendTest() {

@BeforeAll
public static void initialize() throws Exception {
StepVerifier.setDefaultTimeout(Duration.ofSeconds(30));

proxyServer = new SimpleProxy(PROXY_PORT);
proxyServer.start(t -> {
});
Expand All @@ -60,6 +63,8 @@ public void connectFailed(URI uri, SocketAddress sa, IOException ioe) {

@AfterAll
public static void cleanupClient() throws Exception {
StepVerifier.resetDefaultTimeout();

if (proxyServer != null) {
proxyServer.stop();
}
Expand Down Expand Up @@ -89,7 +94,11 @@ public void sendEvents() {
final SendOptions options = new SendOptions().setPartitionId(PARTITION_ID);
final EventHubProducerAsyncClient producer = builder.buildAsyncProducerClient();
final Flux<EventData> events = TestUtils.getEvents(NUMBER_OF_EVENTS, messageId);
final Instant sendTime = Instant.now();
final PartitionProperties information = producer.getPartitionProperties(PARTITION_ID).block();

Assertions.assertNotNull(information, "Should receive partition information.");

final EventPosition position = EventPosition.fromSequenceNumber(information.getLastEnqueuedSequenceNumber());
final EventHubConsumerAsyncClient consumer = builder
.consumerGroup(EventHubClientBuilder.DEFAULT_CONSUMER_GROUP_NAME)
.buildAsyncConsumerClient();
Expand All @@ -101,7 +110,7 @@ public void sendEvents() {
.verify(TIMEOUT);

// Assert
StepVerifier.create(consumer.receiveFromPartition(PARTITION_ID, EventPosition.fromEnqueuedTime(sendTime))
StepVerifier.create(consumer.receiveFromPartition(PARTITION_ID, position)
.filter(x -> TestUtils.isMatchingEvent(x, messageId)).take(NUMBER_OF_EVENTS))
.expectNextCount(NUMBER_OF_EVENTS)
.expectComplete()
Expand Down
Loading