Sonar: Improve test coverage for recent commits
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<String, String> producer = spy(new MockProducer<>());
|
||||
producer.initTransactions();
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
ProducerFactory<String, String> pf = mock(ProducerFactory.class);
|
||||
given(pf.transactionCapable()).willReturn(true);
|
||||
given(pf.createProducer()).willReturn(producer);
|
||||
|
||||
KafkaTemplate<String, String> 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<String, String> producer = spy(new MockProducer<>());
|
||||
producer.initTransactions();
|
||||
producer.fenceProducer();
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
ProducerFactory<String, String> pf = mock(ProducerFactory.class);
|
||||
given(pf.transactionCapable()).willReturn(true);
|
||||
given(pf.createProducer()).willReturn(producer);
|
||||
|
||||
KafkaTemplate<String, String> 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<String, String> producer = spy(new MockProducer<>());
|
||||
producer.initTransactions();
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
ProducerFactory<String, String> pf = mock(ProducerFactory.class);
|
||||
given(pf.transactionCapable()).willReturn(true);
|
||||
given(pf.createProducer()).willReturn(producer);
|
||||
|
||||
KafkaTemplate<String, String> 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 {
|
||||
|
||||
@@ -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<String> ehd = new ErrorHandlingDeserializer<>(new StringDeserializer());
|
||||
assertThat(ehd.deserialize("topic", "foo".getBytes())).isEqualTo("foo");
|
||||
ehd.close();
|
||||
ehd = new ErrorHandlingDeserializer<>(new Deserializer<String>() {
|
||||
|
||||
@Override
|
||||
public void configure(Map<String, ?> 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 {
|
||||
|
||||
Reference in New Issue
Block a user