Test cleanup: Consumer ACK tests

This commit is contained in:
Soby Chacko
2023-02-01 12:04:49 -05:00
parent a7b0cf9d8b
commit cc379b16f9

View File

@@ -31,11 +31,8 @@ import static org.mockito.Mockito.verify;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
@@ -64,11 +61,9 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
void testRecordAck() throws Exception {
Map<String, Object> config = new HashMap<>();
Set<String> strings = new HashSet<>();
strings.add("cons-ack-tests-011");
config.put("topicNames", strings);
config.put("subscriptionName", "cons-ack-tests-sb-011");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("cons-ack-tests-011"),
"subscriptionName", "cons-ack-tests-sb-011");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
@@ -90,8 +85,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
return invocation.callRealMethod();
}).when(containerConsumer).acknowledge(any(MessageId.class));
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "cons-ack-tests-011");
Map<String, Object> prodConfig = Map.of("topicName", "cons-ack-tests-011");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -105,11 +99,8 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
void testBatchAck() throws Exception {
Map<String, Object> config = new HashMap<>();
Set<String> strings = new HashSet<>();
strings.add("cons-ack-tests-012");
config.put("topicNames", strings);
config.put("subscriptionName", "cons-ack-tests-sb-012");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("cons-ack-tests-012"),
"subscriptionName", "cons-ack-tests-sb-012");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
@@ -124,8 +115,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
pulsarConsumerFactory, pulsarContainerProperties);
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "cons-ack-tests-012");
Map<String, Object> prodConfig = Map.of("topicName", "cons-ack-tests-012");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -143,9 +133,8 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
void testBatchAckButSomeRecordsFail() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Collections.singleton("cons-ack-tests-013"));
config.put("subscriptionName", "cons-ack-tests-sb-013");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("cons-ack-tests-013"),
"subscriptionName", "cons-ack-tests-sb-013");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
@@ -171,8 +160,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
return invocation.callRealMethod();
}).when(containerConsumer).acknowledge(any(MessageId.class));
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "cons-ack-tests-013");
Map<String, Object> prodConfig = Map.of("topicName", "cons-ack-tests-013");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -209,11 +197,8 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
@SuppressWarnings("unchecked")
void testManualAckForRecordListener() throws Exception {
Map<String, Object> config = new HashMap<>();
Set<String> strings = new HashSet<>();
strings.add("cons-ack-tests-014");
config.put("topicNames", strings);
config.put("subscriptionName", "cons-ack-tests-sb-014");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("cons-ack-tests-014"),
"subscriptionName", "cons-ack-tests-sb-014");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
@@ -241,8 +226,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
return invocation.callRealMethod();
}).when(containerConsumer).acknowledge(any(MessageId.class));
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "cons-ack-tests-014");
Map<String, Object> prodConfig = Map.of("topicName", "cons-ack-tests-014");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -263,11 +247,8 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
@SuppressWarnings("unchecked")
void testBatchAckForBatchListener() throws Exception {
Map<String, Object> config = new HashMap<>();
Set<String> strings = new HashSet<>();
strings.add("cons-ack-tests-015");
config.put("topicNames", strings);
config.put("subscriptionName", "cons-ack-tests-sb-015");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("cons-ack-tests-015"),
"subscriptionName", "cons-ack-tests-sb-015");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
@@ -291,8 +272,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
pulsarConsumerFactory, pulsarContainerProperties);
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "cons-ack-tests-015");
Map<String, Object> prodConfig = Map.of("topicName", "cons-ack-tests-015");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -311,11 +291,8 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
@SuppressWarnings("unchecked")
void testBatchNackForEntireBatchWhenUsingBatchListener() throws Exception {
Map<String, Object> config = new HashMap<>();
Set<String> strings = new HashSet<>();
strings.add("cons-ack-tests-016");
config.put("topicNames", strings);
config.put("subscriptionName", "cons-ack-tests-sb-016");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("cons-ack-tests-016"),
"subscriptionName", "cons-ack-tests-sb-016");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = spy(
@@ -339,8 +316,7 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
pulsarConsumerFactory, pulsarContainerProperties);
Consumer<String> containerConsumer = ConsumerTestUtils.startContainerAndSpyOnConsumer(container);
Map<String, Object> prodConfig = new HashMap<>();
prodConfig.put("topicName", "cons-ack-tests-016");
Map<String, Object> prodConfig = Map.of("topicName", "cons-ack-tests-016");
DefaultPulsarProducerFactory<String> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient,
prodConfig);
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
@@ -359,9 +335,8 @@ class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
@Test
void messagesAreProperlyAckdOnContainerStopBeforeExitingListenerThread() throws Exception {
Map<String, Object> config = new HashMap<>();
config.put("topicNames", Set.of("duplicate-message-test"));
config.put("subscriptionName", "duplicate-sub-1");
Map<String, Object> config = Map.of("topicNames", Collections.singleton("duplicate-message-test"),
"subscriptionName", "duplicate-sub-1");
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
.build();
DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient,