diff --git a/pom.xml b/pom.xml index 3dd43398..cfc0b940 100644 --- a/pom.xml +++ b/pom.xml @@ -146,6 +146,11 @@ kafka-streams-test-utils test + + org.springframework.kafka + spring-kafka-test + test + dev.vality testcontainers-annotations diff --git a/src/test/java/dev/vality/magista/config/KafkaPostgresqlSpringBootITest.java b/src/test/java/dev/vality/magista/config/KafkaPostgresqlSpringBootITest.java index 97129cd9..c08ee559 100644 --- a/src/test/java/dev/vality/magista/config/KafkaPostgresqlSpringBootITest.java +++ b/src/test/java/dev/vality/magista/config/KafkaPostgresqlSpringBootITest.java @@ -1,8 +1,10 @@ package dev.vality.magista.config; -import dev.vality.testcontainers.annotations.KafkaSpringBootTest; -import dev.vality.testcontainers.annotations.kafka.KafkaTestcontainerSingleton; +import dev.vality.testcontainers.annotations.DefaultSpringBootTest; import dev.vality.testcontainers.annotations.postgresql.PostgresqlTestcontainerSingleton; +import org.springframework.kafka.test.context.EmbeddedKafka; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.TestPropertySource; import java.lang.annotation.ElementType; import java.lang.annotation.Retention; @@ -12,14 +14,20 @@ @Target({ElementType.TYPE}) @Retention(RetentionPolicy.RUNTIME) @PostgresqlTestcontainerSingleton -@KafkaTestcontainerSingleton( - properties = { - "kafka.topics.invoicing.consume.enabled=true", - "kafka.topics.invoice-template.consume.enabled=true", - "kafka.state.cache.size=0"}, - topicsKeys = { - "kafka.topics.invoicing.id", - "kafka.topics.invoice-template.id"}) -@KafkaSpringBootTest +@DefaultSpringBootTest +@EmbeddedKafka(partitions = 1, topics = { + "magista-invoicing-test", + "magista-invoice-template-test" +}) +@TestPropertySource(properties = { + "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}", + "spring.kafka.consumer.group-id=magista-kafka-test", + "kafka.topics.invoicing.id=magista-invoicing-test", + "kafka.topics.invoicing.consume.enabled=true", + "kafka.topics.invoice-template.id=magista-invoice-template-test", + "kafka.topics.invoice-template.consume.enabled=true", + "kafka.state.cache.size=0" +}) +@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) public @interface KafkaPostgresqlSpringBootITest { } diff --git a/src/test/java/dev/vality/magista/kafka/InvoiceTemplateListenerTest.java b/src/test/java/dev/vality/magista/kafka/InvoiceTemplateListenerTest.java index 7f6d50f2..a9bf088f 100644 --- a/src/test/java/dev/vality/magista/kafka/InvoiceTemplateListenerTest.java +++ b/src/test/java/dev/vality/magista/kafka/InvoiceTemplateListenerTest.java @@ -6,13 +6,14 @@ import dev.vality.magista.config.KafkaPostgresqlSpringBootITest; import dev.vality.magista.converter.SourceEventsParser; import dev.vality.magista.service.HandlerManager; -import dev.vality.testcontainers.annotations.kafka.config.KafkaProducer; -import org.apache.thrift.TBase; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; import org.mockito.Captor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.kafka.config.KafkaListenerEndpointRegistry; +import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.test.context.bean.override.mockito.MockitoBean; import java.time.LocalDateTime; @@ -36,13 +37,21 @@ public class InvoiceTemplateListenerTest { private SourceEventsParser sourceEventsParser; @Autowired - private KafkaProducer> testThriftKafkaProducer; + private EmbeddedKafkaBroker embeddedKafkaBroker; + + @Autowired + private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry; @Captor private ArgumentCaptor arg; + @BeforeEach + void waitForKafkaListenersAssignment() { + KafkaTestSupport.waitForAssignments(kafkaListenerEndpointRegistry, embeddedKafkaBroker); + } + @Test - public void shouldInvoiceTemplateSinkEventListen() { + public void shouldInvoiceTemplateSinkEventListen() throws Exception { var message = new MachineEvent(); message.setCreatedAt(LocalDateTime.now().format(DateTimeFormatter.ISO_DATE_TIME)); message.setEventId(1L); @@ -55,7 +64,7 @@ public void shouldInvoiceTemplateSinkEventListen() { sinkEvent.setEvent(message); when(sourceEventsParser.parseEvents(any())) .thenReturn(List.of(EventPayload.invoice_template_changes(List.of()))); - testThriftKafkaProducer.send(invoiceTemplateTopicName, sinkEvent); + KafkaTestSupport.send(embeddedKafkaBroker, invoiceTemplateTopicName, sinkEvent); verify(sourceEventsParser, timeout(5000).times(1)).parseEvents(arg.capture()); assertThat(arg.getValue()) .isEqualTo(message); diff --git a/src/test/java/dev/vality/magista/kafka/InvoicingListenerTest.java b/src/test/java/dev/vality/magista/kafka/InvoicingListenerTest.java index d0a54a2a..e3fe20a6 100644 --- a/src/test/java/dev/vality/magista/kafka/InvoicingListenerTest.java +++ b/src/test/java/dev/vality/magista/kafka/InvoicingListenerTest.java @@ -6,13 +6,14 @@ import dev.vality.magista.config.KafkaPostgresqlSpringBootITest; import dev.vality.magista.converter.SourceEventParser; import dev.vality.magista.service.HandlerManager; -import dev.vality.testcontainers.annotations.kafka.config.KafkaProducer; -import org.apache.thrift.TBase; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; import org.mockito.Captor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.kafka.config.KafkaListenerEndpointRegistry; +import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.test.context.bean.override.mockito.MockitoBean; import java.time.LocalDateTime; @@ -36,13 +37,21 @@ public class InvoicingListenerTest { private SourceEventParser eventParser; @Autowired - private KafkaProducer> testThriftKafkaProducer; + private EmbeddedKafkaBroker embeddedKafkaBroker; + + @Autowired + private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry; @Captor private ArgumentCaptor arg; + @BeforeEach + void waitForKafkaListenersAssignment() { + KafkaTestSupport.waitForAssignments(kafkaListenerEndpointRegistry, embeddedKafkaBroker); + } + @Test - public void shouldInvoicingSinkEventListen() { + public void shouldInvoicingSinkEventListen() throws Exception { var message = new MachineEvent(); message.setCreatedAt(LocalDateTime.now().format(DateTimeFormatter.ISO_DATE_TIME)); message.setEventId(1L); @@ -54,7 +63,7 @@ public void shouldInvoicingSinkEventListen() { var sinkEvent = new SinkEvent(); sinkEvent.setEvent(message); when(eventParser.parseEvent(any())).thenReturn(EventPayload.invoice_changes(List.of())); - testThriftKafkaProducer.send(invoicingTopicName, sinkEvent); + KafkaTestSupport.send(embeddedKafkaBroker, invoicingTopicName, sinkEvent); verify(eventParser, timeout(5000).times(1)).parseEvent(arg.capture()); assertThat(arg.getValue()) .isEqualTo(message); diff --git a/src/test/java/dev/vality/magista/kafka/KafkaTestSupport.java b/src/test/java/dev/vality/magista/kafka/KafkaTestSupport.java new file mode 100644 index 00000000..7ccf172c --- /dev/null +++ b/src/test/java/dev/vality/magista/kafka/KafkaTestSupport.java @@ -0,0 +1,47 @@ +package dev.vality.magista.kafka; + +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.ByteArraySerializer; +import org.apache.kafka.common.serialization.StringSerializer; +import org.apache.thrift.TBase; +import org.apache.thrift.TSerializer; +import org.apache.thrift.protocol.TBinaryProtocol; +import org.springframework.kafka.config.KafkaListenerEndpointRegistry; +import org.springframework.kafka.listener.MessageListenerContainer; +import org.springframework.kafka.test.EmbeddedKafkaBroker; +import org.springframework.kafka.test.utils.ContainerTestUtils; + +import java.util.Properties; + +final class KafkaTestSupport { + + private KafkaTestSupport() { + } + + static void waitForAssignments( + KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry, + EmbeddedKafkaBroker embeddedKafkaBroker) { + for (MessageListenerContainer listenerContainer : kafkaListenerEndpointRegistry.getListenerContainers()) { + ContainerTestUtils.waitForAssignment(listenerContainer, embeddedKafkaBroker.getPartitionsPerTopic()); + } + } + + static void send(EmbeddedKafkaBroker embeddedKafkaBroker, String topic, TBase event) throws Exception { + try (var producer = new KafkaProducer(producerProperties(embeddedKafkaBroker))) { + var serializer = new TSerializer(new TBinaryProtocol.Factory()); + producer.send(new ProducerRecord<>(topic, serializer.serialize(event))).get(); + producer.flush(); + } + } + + private static Properties producerProperties(EmbeddedKafkaBroker embeddedKafkaBroker) { + var properties = new Properties(); + properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, embeddedKafkaBroker.getBrokersAsString()); + properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class); + properties.put(ProducerConfig.ACKS_CONFIG, "all"); + return properties; + } +}