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;
+ }
+}