From 3e45bf84940b1e94c642bf3a88e94514e5a83c41 Mon Sep 17 00:00:00 2001 From: Paul-Christian Volkmer Date: Thu, 29 Feb 2024 09:19:32 +0100 Subject: [PATCH 1/2] feat: implement KafkaInputListener --- .../etl/processor/config/AppConfigurationTest.kt | 2 ++ .../processor/config/AppKafkaConfiguration.kt | 5 +++-- .../etl/processor/input/KafkaInputListener.kt | 16 +++++++++++++--- 3 files changed, 18 insertions(+), 5 deletions(-) diff --git a/src/integrationTest/kotlin/dev/dnpm/etl/processor/config/AppConfigurationTest.kt b/src/integrationTest/kotlin/dev/dnpm/etl/processor/config/AppConfigurationTest.kt index 46a756b..262aca0 100644 --- a/src/integrationTest/kotlin/dev/dnpm/etl/processor/config/AppConfigurationTest.kt +++ b/src/integrationTest/kotlin/dev/dnpm/etl/processor/config/AppConfigurationTest.kt @@ -26,6 +26,7 @@ import dev.dnpm.etl.processor.output.KafkaMtbFileSender import dev.dnpm.etl.processor.output.RestMtbFileSender import dev.dnpm.etl.processor.pseudonym.AnonymizingGenerator import dev.dnpm.etl.processor.pseudonym.GpasPseudonymGenerator +import dev.dnpm.etl.processor.services.RequestProcessor import dev.dnpm.etl.processor.services.TokenRepository import dev.dnpm.etl.processor.services.TokenService import org.assertj.core.api.Assertions.assertThat @@ -144,6 +145,7 @@ class AppConfigurationTest { "app.kafka.group-id=test" ] ) + @MockBean(RequestProcessor::class) inner class AppConfigurationUsingKafkaInputTest(private val context: ApplicationContext) { @Test diff --git a/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt b/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt index 0977527..1ff5e58 100644 --- a/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt +++ b/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt @@ -25,6 +25,7 @@ import dev.dnpm.etl.processor.monitoring.ConnectionCheckService import dev.dnpm.etl.processor.monitoring.KafkaConnectionCheckService import dev.dnpm.etl.processor.output.KafkaMtbFileSender import dev.dnpm.etl.processor.output.MtbFileSender +import dev.dnpm.etl.processor.services.RequestProcessor import dev.dnpm.etl.processor.services.kafka.KafkaResponseProcessor import org.slf4j.LoggerFactory import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean @@ -97,9 +98,9 @@ class AppKafkaConfiguration { @Bean @ConditionalOnProperty(value = ["app.kafka.input-topic"]) fun kafkaInputListener( - applicationEventPublisher: ApplicationEventPublisher, + requestProcessor: RequestProcessor ): KafkaInputListener { - return KafkaInputListener(applicationEventPublisher) + return KafkaInputListener(requestProcessor) } @Bean diff --git a/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt b/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt index ee6b56e..f3c0ab4 100644 --- a/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt +++ b/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt @@ -19,15 +19,25 @@ package dev.dnpm.etl.processor.input +import de.ukw.ccc.bwhc.dto.Consent import de.ukw.ccc.bwhc.dto.MtbFile +import dev.dnpm.etl.processor.services.RequestProcessor import org.apache.kafka.clients.consumer.ConsumerRecord -import org.springframework.context.ApplicationEventPublisher +import org.slf4j.LoggerFactory import org.springframework.kafka.listener.MessageListener class KafkaInputListener( - private val applicationEventPublisher: ApplicationEventPublisher + private val requestProcessor: RequestProcessor ) : MessageListener { + private val logger = LoggerFactory.getLogger(KafkaInputListener::class.java) + override fun onMessage(data: ConsumerRecord) { - TODO("Not yet implemented") + if (data.value().consent.status == Consent.Status.ACTIVE) { + logger.debug("Accepted MTB File for processing") + requestProcessor.processMtbFile(data.value()) + } else { + logger.debug("Accepted MTB File and process deletion") + requestProcessor.processDeletion(data.value().patient.id) + } } } \ No newline at end of file From 952ad8c0cfc64cf9c5e02f1b4f7fc2466f9f2bb3 Mon Sep 17 00:00:00 2001 From: Paul-Christian Volkmer Date: Thu, 29 Feb 2024 12:49:06 +0100 Subject: [PATCH 2/2] test: add test for incoming kafka message processing --- .../processor/config/AppKafkaConfiguration.kt | 5 +- .../etl/processor/input/KafkaInputListener.kt | 15 ++-- .../processor/input/KafkaInputListenerTest.kt | 79 +++++++++++++++++++ 3 files changed, 91 insertions(+), 8 deletions(-) create mode 100644 src/test/kotlin/dev/dnpm/etl/processor/input/KafkaInputListenerTest.kt diff --git a/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt b/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt index 1ff5e58..3799762 100644 --- a/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt +++ b/src/main/kotlin/dev/dnpm/etl/processor/config/AppKafkaConfiguration.kt @@ -98,9 +98,10 @@ class AppKafkaConfiguration { @Bean @ConditionalOnProperty(value = ["app.kafka.input-topic"]) fun kafkaInputListener( - requestProcessor: RequestProcessor + requestProcessor: RequestProcessor, + objectMapper: ObjectMapper ): KafkaInputListener { - return KafkaInputListener(requestProcessor) + return KafkaInputListener(requestProcessor, objectMapper) } @Bean diff --git a/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt b/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt index f3c0ab4..63bf60a 100644 --- a/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt +++ b/src/main/kotlin/dev/dnpm/etl/processor/input/KafkaInputListener.kt @@ -19,6 +19,7 @@ package dev.dnpm.etl.processor.input +import com.fasterxml.jackson.databind.ObjectMapper import de.ukw.ccc.bwhc.dto.Consent import de.ukw.ccc.bwhc.dto.MtbFile import dev.dnpm.etl.processor.services.RequestProcessor @@ -27,17 +28,19 @@ import org.slf4j.LoggerFactory import org.springframework.kafka.listener.MessageListener class KafkaInputListener( - private val requestProcessor: RequestProcessor -) : MessageListener { + private val requestProcessor: RequestProcessor, + private val objectMapper: ObjectMapper +) : MessageListener { private val logger = LoggerFactory.getLogger(KafkaInputListener::class.java) - override fun onMessage(data: ConsumerRecord) { - if (data.value().consent.status == Consent.Status.ACTIVE) { + override fun onMessage(data: ConsumerRecord) { + val mtbFile = objectMapper.readValue(data.value(), MtbFile::class.java) + if (mtbFile.consent.status == Consent.Status.ACTIVE) { logger.debug("Accepted MTB File for processing") - requestProcessor.processMtbFile(data.value()) + requestProcessor.processMtbFile(mtbFile) } else { logger.debug("Accepted MTB File and process deletion") - requestProcessor.processDeletion(data.value().patient.id) + requestProcessor.processDeletion(mtbFile.patient.id) } } } \ No newline at end of file diff --git a/src/test/kotlin/dev/dnpm/etl/processor/input/KafkaInputListenerTest.kt b/src/test/kotlin/dev/dnpm/etl/processor/input/KafkaInputListenerTest.kt new file mode 100644 index 0000000..cf5ba39 --- /dev/null +++ b/src/test/kotlin/dev/dnpm/etl/processor/input/KafkaInputListenerTest.kt @@ -0,0 +1,79 @@ +/* + * This file is part of ETL-Processor + * + * Copyright (c) 2024 Comprehensive Cancer Center Mainfranken, Datenintegrationszentrum Philipps-Universität Marburg and Contributors + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +package dev.dnpm.etl.processor.input + +import com.fasterxml.jackson.databind.ObjectMapper +import de.ukw.ccc.bwhc.dto.Consent +import de.ukw.ccc.bwhc.dto.MtbFile +import de.ukw.ccc.bwhc.dto.Patient +import dev.dnpm.etl.processor.services.RequestProcessor +import org.apache.kafka.clients.consumer.ConsumerRecord +import org.junit.jupiter.api.BeforeEach +import org.junit.jupiter.api.Test +import org.junit.jupiter.api.extension.ExtendWith +import org.mockito.ArgumentMatchers.anyString +import org.mockito.Mock +import org.mockito.junit.jupiter.MockitoExtension +import org.mockito.kotlin.any +import org.mockito.kotlin.times +import org.mockito.kotlin.verify + +@ExtendWith(MockitoExtension::class) +class KafkaInputListenerTest { + + private lateinit var requestProcessor: RequestProcessor + private lateinit var objectMapper: ObjectMapper + private lateinit var kafkaInputListener: KafkaInputListener + + @BeforeEach + fun setup( + @Mock requestProcessor: RequestProcessor + ) { + this.requestProcessor = requestProcessor + this.objectMapper = ObjectMapper() + + this.kafkaInputListener = KafkaInputListener(requestProcessor, objectMapper) + } + + @Test + fun shouldProcessMtbFileRequest() { + val mtbFile = MtbFile.builder() + .withPatient(Patient.builder().withId("DUMMY_12345678").build()) + .withConsent(Consent.builder().withStatus(Consent.Status.ACTIVE).build()) + .build() + + kafkaInputListener.onMessage(ConsumerRecord("testtopic", 0, 0, "", this.objectMapper.writeValueAsString(mtbFile))) + + verify(requestProcessor, times(1)).processMtbFile(any()) + } + + @Test + fun shouldProcessDeleteRequest() { + val mtbFile = MtbFile.builder() + .withPatient(Patient.builder().withId("DUMMY_12345678").build()) + .withConsent(Consent.builder().withStatus(Consent.Status.REJECTED).build()) + .build() + + kafkaInputListener.onMessage(ConsumerRecord("testtopic", 0, 0, "", this.objectMapper.writeValueAsString(mtbFile))) + + verify(requestProcessor, times(1)).processDeletion(anyString()) + } + +} \ No newline at end of file