From 5ce42d2e42ee06c6610c869d12fc3dc76903f7f3 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Sun, 4 Nov 2018 14:31:39 -0500 Subject: [PATCH] Sonar: Improve test coverage for recent commits --- .../serializer/DeserializationException.java | 2 +- .../core/KafkaTemplateTransactionTests.java | 74 +++++++++++++++++++ .../ErrorHandlingDeserializerTests.java | 28 +++++++ 3 files changed, 103 insertions(+), 1 deletion(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/DeserializationException.java b/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/DeserializationException.java index 3192d4dd..91fb4587 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/DeserializationException.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/serializer/DeserializationException.java @@ -42,7 +42,7 @@ public class DeserializationException extends KafkaException { private final boolean isKey; public DeserializationException(String message, byte[] data, boolean isKey, Throwable cause) { - this(message, null, data, isKey, cause); + this(message, null, data, isKey, cause); // NOSONAR test coverage } public DeserializationException(String message, @Nullable Headers headers, byte[] data, // NOSONAR array reference diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java index 7131d3ae..92781c01 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java @@ -23,6 +23,8 @@ import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.springframework.kafka.test.assertj.KafkaConditions.key; @@ -41,6 +43,7 @@ import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.errors.ProducerFencedException; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.assertj.core.api.Assertions; @@ -266,6 +269,77 @@ public class KafkaTemplateTransactionTests { assertThat(producer.closed()).isTrue(); } + @Test + public void testNoAbortAfterCommitFailure() { + MockProducer producer = spy(new MockProducer<>()); + producer.initTransactions(); + + @SuppressWarnings("unchecked") + ProducerFactory pf = mock(ProducerFactory.class); + given(pf.transactionCapable()).willReturn(true); + given(pf.createProducer()).willReturn(producer); + + KafkaTemplate template = new KafkaTemplate<>(pf); + template.setDefaultTopic(STRING_KEY_TOPIC); + + assertThatThrownBy(() -> template.executeInTransaction(t -> { + producer.fenceProducer(); + return null; + })).isInstanceOf(ProducerFencedException.class); + + assertThat(producer.transactionCommitted()).isFalse(); + assertThat(producer.transactionAborted()).isFalse(); + assertThat(producer.closed()).isTrue(); + verify(producer, never()).abortTransaction(); + } + + @Test + public void testFencedOnBegin() { + MockProducer producer = spy(new MockProducer<>()); + producer.initTransactions(); + producer.fenceProducer(); + + @SuppressWarnings("unchecked") + ProducerFactory pf = mock(ProducerFactory.class); + given(pf.transactionCapable()).willReturn(true); + given(pf.createProducer()).willReturn(producer); + + KafkaTemplate template = new KafkaTemplate<>(pf); + template.setDefaultTopic(STRING_KEY_TOPIC); + + assertThatThrownBy(() -> template.executeInTransaction(t -> { + return null; + })).isInstanceOf(ProducerFencedException.class); + + assertThat(producer.transactionCommitted()).isFalse(); + assertThat(producer.transactionAborted()).isFalse(); + assertThat(producer.closed()).isTrue(); + verify(producer, never()).commitTransaction(); + } + + @Test + public void testAbort() { + MockProducer producer = spy(new MockProducer<>()); + producer.initTransactions(); + + @SuppressWarnings("unchecked") + ProducerFactory pf = mock(ProducerFactory.class); + given(pf.transactionCapable()).willReturn(true); + given(pf.createProducer()).willReturn(producer); + + KafkaTemplate template = new KafkaTemplate<>(pf); + template.setDefaultTopic(STRING_KEY_TOPIC); + + assertThatThrownBy(() -> template.executeInTransaction(t -> { + throw new RuntimeException("foo"); + })).isExactlyInstanceOf(RuntimeException.class).withFailMessage("foo"); + + assertThat(producer.transactionCommitted()).isFalse(); + assertThat(producer.transactionAborted()).isTrue(); + assertThat(producer.closed()).isTrue(); + verify(producer, never()).commitTransaction(); + } + @Configuration @EnableTransactionManagement public static class DeclarativeConfig { diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ErrorHandlingDeserializerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ErrorHandlingDeserializerTests.java index ccafef9c..fb85dc63 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ErrorHandlingDeserializerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ErrorHandlingDeserializerTests.java @@ -26,7 +26,9 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.header.Headers; +import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.ExtendedDeserializer; +import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.junit.jupiter.api.Test; @@ -74,6 +76,32 @@ public class ErrorHandlingDeserializerTests { assertThat(this.config.headers).isNotNull(); } + @Test + public void unitTests() { + ErrorHandlingDeserializer ehd = new ErrorHandlingDeserializer<>(new StringDeserializer()); + assertThat(ehd.deserialize("topic", "foo".getBytes())).isEqualTo("foo"); + ehd.close(); + ehd = new ErrorHandlingDeserializer<>(new Deserializer() { + + @Override + public void configure(Map configs, boolean isKey) { + } + + @Override + public String deserialize(String topic, byte[] data) { + throw new RuntimeException("fail"); + } + + @Override + public void close() { + } + + }); + Object result = ehd.deserialize("topic", "foo".getBytes()); + assertThat(result).isInstanceOf(DeserializationException.class); + ehd.close(); + } + @Configuration @EnableKafka public static class Config {