Enhance PulsarConsumerTestUtil consumeMessages

This commit allows users to pass in no args, pulsar client,
or broker url to PulsarConsumerTestUtil consumeMessages().

Resolves #599
This commit is contained in:
KartikShrivastava
2024-03-10 21:19:53 +05:30
committed by Chris Bono
parent b55f5bff57
commit 2c7e3b248a
3 changed files with 136 additions and 2 deletions

View File

@@ -24,11 +24,13 @@ import java.util.concurrent.TimeUnit;
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.SubscriptionInitialPosition;
import org.springframework.pulsar.PulsarException;
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.util.Assert;
@@ -53,6 +55,30 @@ public class PulsarConsumerTestUtil<T> implements TopicSpec<T>, SchemaSpec<T>, C
private List<String> topics;
private boolean untilMethodAlreadyCalled = false;
public static <T> TopicSpec<T> consumeMessages() {
if (PulsarTestContainerSupport.isContainerStarted()) {
return PulsarConsumerTestUtil.consumeMessages(PulsarTestContainerSupport.getPulsarBrokerUrl());
}
return PulsarConsumerTestUtil.consumeMessages("pulsar://localhost:6650");
}
public static <T> TopicSpec<T> consumeMessages(String url) {
Assert.notNull(url, "url must not be null");
try {
return PulsarConsumerTestUtil.consumeMessages(PulsarClient.builder().serviceUrl(url).build());
}
catch (PulsarClientException ex) {
throw new PulsarException(ex);
}
}
public static <T> TopicSpec<T> consumeMessages(PulsarClient pulsarClient) {
Assert.notNull(pulsarClient, "pulsarClient must not be null");
return PulsarConsumerTestUtil.consumeMessages(new DefaultPulsarConsumerFactory<>(pulsarClient, List.of()));
}
public static <T> TopicSpec<T> consumeMessages(PulsarConsumerFactory<T> pulsarConsumerFactory) {
return new PulsarConsumerTestUtil<>(pulsarConsumerFactory);
}
@@ -85,6 +111,11 @@ public class PulsarConsumerTestUtil<T> implements TopicSpec<T>, SchemaSpec<T>, C
@Override
public ConditionsSpec<T> until(ConsumedMessagesCondition<T> condition) {
if (untilMethodAlreadyCalled) {
throw new IllegalStateException(
"Multiple calls to 'until' are not allowed. Use 'and' to combine conditions.");
}
this.untilMethodAlreadyCalled = true;
this.condition = condition;
return this;
}

View File

@@ -48,4 +48,12 @@ public interface PulsarTestContainerSupport {
return PULSAR_CONTAINER.getHttpServiceUrl();
}
static boolean isContainerStarted() {
return PULSAR_CONTAINER.isRunning();
}
static void stopContainer() {
PULSAR_CONTAINER.stop();
}
}

View File

@@ -34,6 +34,7 @@ import org.junit.jupiter.api.Test;
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.pulsar.core.PulsarTemplate;
/**
@@ -47,13 +48,15 @@ class PulsarConsumerTestUtilTests implements PulsarTestContainerSupport {
private DefaultPulsarConsumerFactory<String> pulsarConsumerFactory;
private PulsarClient pulsarClient;
private static String testTopic(String suffix) {
return "ptctut-topic-" + suffix;
}
@BeforeEach
void prepareForTest() throws PulsarClientException {
var pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
this.pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
this.pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, List.of());
this.pulsarTemplate = new PulsarTemplate<>(new DefaultPulsarProducerFactory<>(pulsarClient));
}
@@ -97,9 +100,101 @@ class PulsarConsumerTestUtilTests implements PulsarTestContainerSupport {
.withMessage("Condition was not met within 5 seconds");
}
@Test
void messagesAreConsumedWhenContainerIsRunningAndConsumeMessagesIsCalledWithoutArguments() {
// depends upon pulsarClient created in prepareForTest
var topic = testTopic("e1");
IntStream.range(0, 2).forEach(i -> pulsarTemplate.send(topic, "message-" + i));
var msgs = PulsarConsumerTestUtil.<String>consumeMessages()
.fromTopic(topic)
.withSchema(Schema.STRING)
.awaitAtMost(Duration.ofSeconds(5))
.until(desiredMessageCount(2))
.get();
assertThat(msgs).hasSize(2);
}
@Test
void messagesAreConsumedWhenContainerIsStoppedAndConsumeMessagesIsCalledWithoutArguments() {
PulsarTestContainerSupport.stopContainer();
PulsarConsumerTestUtil.<String>consumeMessages();
// TODO: Complete this test
}
@Test
void messagesAreConsumedWhenConsumeMessagesIsCalledWithBrokerUrl() {
var topic = testTopic("e2");
IntStream.range(0, 2).forEach(i -> pulsarTemplate.send(topic, "message-" + i));
var msgs = PulsarConsumerTestUtil.<String>consumeMessages(PulsarTestContainerSupport.getPulsarBrokerUrl())
.fromTopic(topic)
.withSchema(Schema.STRING)
.awaitAtMost(Duration.ofSeconds(5))
.until(desiredMessageCount(2))
.get();
assertThat(msgs).hasSize(2);
}
@Test
void exceptionIsThrownWhenConsumeMessagesIsCalledWithNullBrokerUrl() {
String url = null;
assertThatIllegalArgumentException().isThrownBy(() -> PulsarConsumerTestUtil.consumeMessages(url))
.withMessage("url must not be null");
}
@Test
void messagesAreConsumedWhenConsumeMessagesIsCalledWithPulsarClient() {
var topic = testTopic("e3");
IntStream.range(0, 2).forEach(i -> pulsarTemplate.send(topic, "message-" + i));
var msgs = PulsarConsumerTestUtil.<String>consumeMessages(this.pulsarClient)
.fromTopic(topic)
.withSchema(Schema.STRING)
.awaitAtMost(Duration.ofSeconds(5))
.until(desiredMessageCount(2))
.get();
assertThat(msgs).hasSize(2);
}
@Test
void exceptionIsThrownWhenConsumeMessagesIsCalledWithNullPulsarClient() {
PulsarClient localPulsarClient = null;
assertThatIllegalArgumentException().isThrownBy(() -> PulsarConsumerTestUtil.consumeMessages(localPulsarClient))
.withMessage("pulsarClient must not be null");
}
@Test
void whenChainedConditionAreSpecifiedMessagesAreConsumedUntilTheyAreMet() {
var topic = testTopic("d");
IntStream.range(0, 5).forEach(i -> pulsarTemplate.send(topic, "message-" + i));
ConsumedMessagesCondition<String> condition1 = ConsumedMessagesConditions.desiredMessageCount(5);
ConsumedMessagesCondition<String> condition2 = ConsumedMessagesConditions.atLeastOneMessageMatches("message-1");
var msgs = PulsarConsumerTestUtil.consumeMessages(pulsarConsumerFactory)
.fromTopic(topic)
.withSchema(Schema.STRING)
.awaitAtMost(Duration.ofSeconds(5))
.until(condition1.and(condition2))
.get();
assertThat(msgs).hasSize(5);
}
@Test
void exceptionIsThrownWhenUntilIsCalledMultipleTimes() {
var topic = testTopic("e");
IntStream.range(0, 1).forEach(i -> pulsarTemplate.send(topic, "message-" + i));
assertThatExceptionOfType(IllegalStateException.class)
.isThrownBy(() -> PulsarConsumerTestUtil.consumeMessages(pulsarConsumerFactory)
.fromTopic(topic)
.withSchema(Schema.STRING)
.awaitAtMost(Duration.ofSeconds(5))
.until(ConsumedMessagesConditions.desiredMessageCount(1))
.until(ConsumedMessagesConditions.atLeastOneMessageMatches("message-0"))
.get())
.withMessage("Multiple calls to 'until' are not allowed. Use 'and' to combine conditions.");
}
@Test
void consumerFactoryCannotBeNull() {
assertThatIllegalArgumentException().isThrownBy(() -> PulsarConsumerTestUtil.consumeMessages(null))
PulsarConsumerFactory<String> consumerFactory = null;
assertThatIllegalArgumentException().isThrownBy(() -> PulsarConsumerTestUtil.consumeMessages(consumerFactory))
.withMessage("PulsarConsumerFactory must not be null");
}