From a915121c3f4871c5841a21e43f5341981afd4662 Mon Sep 17 00:00:00 2001 From: psemenkov Date: Thu, 2 Feb 2023 11:36:22 +0300 Subject: [PATCH] Add Mock kafka producer. Fixing test in balance-service. --- .../queue/AbstractControllerTest.java | 13 +- .../PaymentInstructionControllerTest.java | 4 +- clearing-parent/balance-service/pom.xml | 10 ++ ...mdgTestConfig.java => ImdgTestConfig.java} | 4 +- .../balance/config/KafkaTestConfig.java | 12 +- .../balance/service/AbstractServiceTest.java | 53 ++++--- .../service/AccountBalanceServiceTest.java | 7 +- .../balance/service/Sdf01ExecutorTest.java | 26 ++-- .../balance/service/Sdf08ServiceTest.java | 62 ++++---- .../balance/service/Sdf09ExecutorTest.java | 4 + .../balance/service/Sdf16ExecutorTest.java | 24 ++++ .../service/StatementServiceServiceTest.java | 132 +++++++++++++----- .../balance/utils/MockKafkaUtils.java | 54 +++++++ .../messaging/service/QueueConsumer.java | 15 +- 14 files changed, 288 insertions(+), 132 deletions(-) rename clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/{BalanceImdgTestConfig.java => ImdgTestConfig.java} (96%) create mode 100644 clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MockKafkaUtils.java diff --git a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java index 4ef207ea7..fbcecd3b0 100644 --- a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java +++ b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java @@ -20,9 +20,9 @@ import ru.spcex.clearing.backendapi.controller.queue.company.*; import ru.spcex.clearing.backendapi.controller.queue.misc.NotificationController; import ru.spcex.clearing.backendapi.controller.queue.misc.SessionController; import ru.spcex.clearing.backendapi.controller.queue.payment.PaymentInstructionController; -import ru.spcex.clearing.backendapi.controller.queue.scheduler.SchedulerAllTodayController; -import ru.spcex.clearing.backendapi.controller.queue.scheduler.SchedulerController; -import ru.spcex.clearing.backendapi.controller.queue.scheduler.TaskRunnerController; +import ru.spcex.clearing.backendapi.controller.queue.scheduler.ClearingCalendarController; +import ru.spcex.clearing.backendapi.controller.queue.scheduler.LauncherController; +import ru.spcex.clearing.backendapi.controller.queue.scheduler.PlannerAllTodayController; import ru.spcex.clearing.backendapi.meta.GetResponseFactory; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; @@ -47,11 +47,12 @@ import java.util.concurrent.atomic.AtomicLong; //payment PaymentInstructionController.class, //scheduler - SchedulerAllTodayController.class, - SchedulerController.class, - TaskRunnerController.class, + ClearingCalendarController.class, + PlannerAllTodayController.class,//PlanerAllTodayController + LauncherController.class, //******* common configs ******* WebTestConfig.class, + GetResponseFactory.class, IOperatorTest.class, KafkaTestConfig.class, HazelcastServiceTestConfiguration.class, diff --git a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java index 3fba79694..e2df876d8 100644 --- a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java +++ b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java @@ -34,9 +34,9 @@ class PaymentInstructionControllerTest extends AbstractControllerTest { PaymentInstruction paymentInstruction = new PaymentInstruction(); paymentInstruction.setSenderId(1010L); paymentInstruction.setAddresseeId(1010L); - paymentInstruction.setAdresseeBIC("123456789"); + paymentInstruction.setAdresseeBic("123456789"); paymentInstruction.setPayeeBankName("BankName"); - paymentInstruction.setPayeeBIC("123456789"); + paymentInstruction.setPayeeBic("123456789"); paymentInstruction.setAddresseeBankName("BankName"); paymentInstruction.setPaymentDate(Instant.now()); paymentInstruction.setPaymentPurpose("Purpose"); diff --git a/clearing-parent/balance-service/pom.xml b/clearing-parent/balance-service/pom.xml index d49af4c57..f4d1145bb 100644 --- a/clearing-parent/balance-service/pom.xml +++ b/clearing-parent/balance-service/pom.xml @@ -56,6 +56,16 @@ assertj-core test + + org.springframework.boot + spring-boot-test + test + + + org.mockito + mockito-core + test + diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/BalanceImdgTestConfig.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/ImdgTestConfig.java similarity index 96% rename from clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/BalanceImdgTestConfig.java rename to clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/ImdgTestConfig.java index 6b8a4e34c..c629b65d7 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/BalanceImdgTestConfig.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/ImdgTestConfig.java @@ -17,7 +17,7 @@ import java.util.List; import java.util.Random; @Configuration -public class BalanceImdgTestConfig { +public class ImdgTestConfig { private HazelcastInstance hazelcastInstance; @@ -57,7 +57,7 @@ public class BalanceImdgTestConfig { joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1"))); networkConfig.setJoin(joinConfig); cfg.setNetworkConfig(networkConfig); - hazelcastInstance = Hazelcast.newHazelcastInstance(cfg); + hazelcastInstance = Hazelcast.getOrCreateHazelcastInstance(cfg); HazelcastHelper.otcSystem_setStorageState(true, hazelcastInstance); return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params); } diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java index b56fef2b1..6a533e178 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/config/KafkaTestConfig.java @@ -2,23 +2,17 @@ package ru.spcex.clearing.balance.config; import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.consumer.OffsetResetStrategy; -import org.apache.kafka.clients.producer.MockProducer; -import org.apache.kafka.clients.producer.Producer; -import org.apache.kafka.common.serialization.StringSerializer; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import ru.spcex.clearing.platform.messaging.serialization.JsonSerializer; +import org.springframework.context.annotation.Scope; @Configuration public class KafkaTestConfig { + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) @Bean public MockConsumer createTestConsumer() { return new MockConsumer<>(OffsetResetStrategy.EARLIEST); } - - @Bean - public Producer createTestProducer() { - return new MockProducer<>(true, new StringSerializer(), new JsonSerializer()); - } } diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java index 1b1aad983..2064efd64 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AbstractServiceTest.java @@ -10,10 +10,7 @@ import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.AccountBalance; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.CompanySymbols; -import ru.clearing.classes.statics.data.sdf.SDf02; -import ru.clearing.classes.statics.data.sdf.SDf08; -import ru.clearing.classes.statics.data.sdf.SDf10; -import ru.clearing.classes.statics.data.sdf.SDf17; +import ru.clearing.classes.statics.data.sdf.*; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.balance.config.*; import ru.spcex.clearing.balance.utils.MatcherFactory; @@ -25,8 +22,11 @@ import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast; import java.math.BigDecimal; import java.math.RoundingMode; +import java.time.LocalDate; import java.util.concurrent.atomic.AtomicLong; +import static ru.spcex.clearing.balance.service.Sdf01ExecutorTest.acc; +import static ru.spcex.clearing.balance.service.Sdf01ExecutorTest.datFormatter; import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFieldsComparator; @ExtendWith(SpringExtension.class) @@ -41,13 +41,27 @@ import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFields MessagesConfig.class, ValidationConfig.class, SdfExecutorsConfig.class, - BalanceImdgTestConfig.class, + ImdgTestConfig.class, KafkaSenderConfig.class, + ImdgTestConfig.class, KafkaTestConfig.class}) public abstract class AbstractServiceTest { protected static final MatcherFactory.Matcher ACCOUNT_BALANCE_MATCHER = usingIgnoringFieldsComparator("created", "updated", "clearingDate", "id"); protected static final MatcherFactory.Matcher RESULT_MATCHER = usingIgnoringFieldsComparator("account.created", "account.updated", "account.clearingDate", "generationId"); protected static final MatcherFactory.Matcher STATEMENT_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "updated", "id"); + protected static AtomicLong currentId = new AtomicLong(0L); + protected static IMap companyMap; + protected static IMap companySymbolsMap; + protected static IMap accountMap; + protected static ImdgHazelcast statementImdg; + protected static ImdgHazelcast requestInfoImdg; + protected static ImdgHazelcast sdf02Imdg; + protected static ImdgHazelcast sdf01Imdg; + protected static ImdgHazelcast sdf08Imdg; + protected static ImdgHazelcast sdf10Imdg; + protected static ImdgHazelcast sdf17Imdg; + protected static ImdgHazelcast accountBalanceImdg; + protected static ImdgHazelcast companyImdg; protected final Long accountIdNew = 10L; protected final Long addresseeIdNew = 2L; protected final String deal = "111111111"; @@ -55,28 +69,18 @@ public abstract class AbstractServiceTest { protected final BigDecimal amountNew = new BigDecimal(1000); protected final String cashMovementCurrencyCode = CurrencyCode.RUB.getKey(); - protected AtomicLong currentId = new AtomicLong(0L); - protected IMap companyMap; - protected IMap companySymbolsMap; - protected IMap accountMap; - protected ImdgHazelcast statementImdg; - protected ImdgHazelcast requestInfoImdg; - protected ImdgHazelcast sdf02Imdg; - protected ImdgHazelcast sdf08Imdg; - protected ImdgHazelcast sdf10Imdg; - protected ImdgHazelcast sdf17Imdg; - protected ImdgHazelcast accountBalanceImdg; @Autowired @Qualifier("hazelcastServiceTest") private ImdgProvider hazelcast; void init() { hazelcast.waitAvailable(); - ImdgHazelcast companyImdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class); + companyImdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class); ImdgHazelcast companySymbolsImdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); ImdgHazelcast accountImdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_Account, Account.class); requestInfoImdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); statementImdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); + sdf01Imdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); sdf02Imdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class); sdf08Imdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class); sdf10Imdg = (ImdgHazelcast) hazelcast.getImdg(IMDGDistributedNames.Map_SDf10, SDf10.class); @@ -136,4 +140,19 @@ public abstract class AbstractServiceTest { accountBalance.setFullName(company.getFullName()); return accountBalance; } + + protected SDf01 getTestSdf01(long id, long generationId) { + LocalDate date = LocalDate.now(); + SDf01 sdf01 = new SDf01(); + sdf01.setId(id); + sdf01.setGenerationId(generationId); + sdf01.setMarket("U"); + sdf01.setDeal(deal); + sdf01.setAccount(acc); + sdf01.setCurr_code("RUR"); + sdf01.setDat(date.format(datFormatter)); + sdf01.setAcc_type("A"); + sdf01.setRemainder(amountNew.toString()); + return sdf01; + } } diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java index 2ae63f1bf..90849e793 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/AccountBalanceServiceTest.java @@ -1,7 +1,9 @@ package ru.spcex.clearing.balance.service; +import org.apache.kafka.clients.producer.MockProducer; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.mock.mockito.SpyBean; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.AccountBalance; import ru.clearing.classes.statics.data.company.Company; @@ -30,11 +32,14 @@ class AccountBalanceServiceTest extends AbstractServiceTest { @Autowired AccountBalanceService accountBalanceService; + @SpyBean + private MockProducer producer; + @PostConstruct void init() { super.init(); } - + /** * {@link AccountBalanceService#createAccountBalance(Long addresseeId, Long accountId, BigDecimal amount, String cashMovementCurrencyCode)}
* Тест проверяет генерацию сущности {@link AccountResult}
diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf01ExecutorTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf01ExecutorTest.java index 4aa813da4..dfd5b6b90 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf01ExecutorTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf01ExecutorTest.java @@ -1,7 +1,9 @@ package ru.spcex.clearing.balance.service; +import org.apache.kafka.clients.producer.MockProducer; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.mock.mockito.SpyBean; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.AccountBalance; import ru.clearing.classes.statics.data.company.Company; @@ -27,13 +29,15 @@ import java.util.Map; import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFieldsComparator; class Sdf01ExecutorTest extends AbstractServiceTest { - final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy"); + public final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy"); + public final static String acc = "123456789"; private static final MatcherFactory.Matcher SDF_02_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id"); - private final String acc = "123456789"; private final Long ID = 1L; @Autowired Sdf01Executor sdf01Executor; private StatementRequest statementRequest; + @SpyBean + private MockProducer producer; @PostConstruct void init() { @@ -59,7 +63,7 @@ class Sdf01ExecutorTest extends AbstractServiceTest { @Test void execute() { //check create statement, SDf02 and AccountBalance - SDf01 sdf01 = getTestSdf01(); + SDf01 sdf01 = getTestSdf01(ID, 1L); Company company = getTestCompany(); companyMap.put(addresseeIdNew, company); @@ -103,7 +107,7 @@ class Sdf01ExecutorTest extends AbstractServiceTest { @Test void validatedExecute() { //CompanyNotFound - SDf01 sdf01 = getTestSdf01(); + SDf01 sdf01 = getTestSdf01(ID, 1L); companyMap.delete(addresseeIdNew); checkError(new EnumMessage(BalanceError.CompanyNotFound), sdf01); @@ -206,20 +210,6 @@ class Sdf01ExecutorTest extends AbstractServiceTest { return sDf02; } - private SDf01 getTestSdf01() { - LocalDate date = LocalDate.now(); - SDf01 sdf01 = new SDf01(); - sdf01.setId(ID); - sdf01.setMarket("U"); - sdf01.setDeal(deal); - sdf01.setAccount(acc); - sdf01.setCurr_code("RUR"); - sdf01.setDat(date.format(datFormatter)); - sdf01.setAcc_type("A"); - sdf01.setRemainder(amountNew.toString()); - return sdf01; - } - private Statement getTestStatement(Long id, Company company, Account account, SDf01 sdf01) { Statement statement = new Statement(); statement.setId(id); diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceTest.java index afb4fd05d..f78043233 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf08ServiceTest.java @@ -2,43 +2,46 @@ package ru.spcex.clearing.balance.service; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -import com.hazelcast.core.IMap; -import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.MockConsumer; -import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.clients.producer.MockProducer; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.mock.mockito.SpyBean; import ru.clearing.classes.statics.data.sdf.SDf08; -import ru.spcex.clearing.balance.utils.ImapEvent; import ru.spcex.clearing.platform.messaging.domain.ActionType; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest; import ru.spcex.clearing.platform.messaging.service.RequestInfo; -import ru.spcex.clearing.platform.messaging.service.Status; import ru.spcex.platform.enumeration.Task; import javax.annotation.PostConstruct; -import java.time.Instant; -import java.util.Collection; -import java.util.Collections; -import java.util.HashMap; -import java.util.Map; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.mockito.Mockito.timeout; +import static org.mockito.Mockito.verify; +import static ru.spcex.clearing.balance.utils.MockKafkaUtils.addRecordToKafka; class Sdf08ServiceTest extends AbstractServiceTest { // private static final MatcherFactory.Matcher SDF_08_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id"); private static final String TOPIC = Task.getAllBalance.topic(); private static final int PARTITION = 1; - private static final Long limit = 300L; @Autowired Sdf08Service sdf08Service; - @Autowired + @Captor + ArgumentCaptor producerRecord; private MockConsumer mockConsumer; + @SpyBean + private MockProducer producer; @PostConstruct void init() { super.init(); + mockConsumer = (MockConsumer) sdf08Service.getConsumer(); } /** @@ -60,30 +63,21 @@ class Sdf08ServiceTest extends AbstractServiceTest { throw new RuntimeException(e); } - IMap sdf08Map = sdf08Imdg.getMap(); - IMap resultsRequestMap = requestInfoImdg.getMap(); - ImapEvent imapEvent = new ImapEvent(resultsRequestMap); - int count = sdf08Map.size(); - long timeNow = Instant.now().toEpochMilli(); - //KAFKA - mockConsumer.schedulePollTask(() -> { - mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC, PARTITION))); - mockConsumer.addRecord(new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonBaseNewRequest)); - }); + addRecordToKafka(mockConsumer, TOPIC, PARTITION, 0, jsonBaseNewRequest); - HashMap startOffsets = new HashMap<>(); - TopicPartition tp = new TopicPartition(TOPIC, PARTITION); - startOffsets.put(tp, 0L); - mockConsumer.updateBeginningOffsets(startOffsets); + //waiting for kafka producer send message (finale event) + verify(producer, timeout(30_000L).times(1)) + .send(producerRecord.capture()); - //waiting for hazelcast map item updates - imapEvent.waitWhenHappened(); + assertEquals(Consts.EXPORT_PROCESS, producerRecord.getValue().topic()); + BaseRequest baseRequest = (BaseRequest) producerRecord.getValue().value(); - Collection resultsRequestInfo = requestInfoImdg.getCollectionObjectsByFieldValues(Map.of("status", Status.Processing)); - RequestInfo requestInfo = resultsRequestInfo.stream().max((entry1, entry2) -> entry1.getId() > entry2.getId() ? 1 : -1).get(); - Long diffRequestInfo = requestInfo.getCreated().toEpochMilli() - timeNow; + ExportToFileRequest exportToFileRequest = (ExportToFileRequest) baseRequest.getRequestPayload(); + SDf08 resultsSDf08 = sdf08Imdg.getSingleObjectBySQL(String.format("generationId = %s and id != null", exportToFileRequest.getSdfGroupId())); + RequestInfo resultRequestInfo = requestInfoImdg.getSingleObjectByID(baseRequest.getId()); - assertEquals(1, sdf08Map.size() - count); - assertTrue(limit.compareTo(diffRequestInfo) > 0); + assertNotNull(baseRequest); + assertNotNull(resultsSDf08); + assertNotNull(resultRequestInfo); } } \ No newline at end of file diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf09ExecutorTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf09ExecutorTest.java index db4a23ebd..8cd52bfa0 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf09ExecutorTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf09ExecutorTest.java @@ -1,7 +1,9 @@ package ru.spcex.clearing.balance.service; +import org.apache.kafka.clients.producer.MockProducer; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.mock.mockito.SpyBean; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.AccountBalance; import ru.clearing.classes.statics.data.company.Company; @@ -34,6 +36,8 @@ class Sdf09ExecutorTest extends AbstractServiceTest { Sdf09Executor sdf09Executor; private SDf09 sdf09; private StatementRequest statementRequest; + @SpyBean + private MockProducer producer; @PostConstruct void init() { diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf16ExecutorTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf16ExecutorTest.java index c475d007d..6761fdd92 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf16ExecutorTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/Sdf16ExecutorTest.java @@ -1,7 +1,12 @@ package ru.spcex.clearing.balance.service; +import org.apache.kafka.clients.producer.MockProducer; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.mock.mockito.MockBean; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.AccountBalance; import ru.clearing.classes.statics.data.company.Company; @@ -11,8 +16,12 @@ import ru.clearing.classes.statics.data.sdf.SDf17; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.balance.errors.BalanceError; import ru.spcex.clearing.balance.utils.MatcherFactory; +import ru.spcex.clearing.balance.utils.MockKafkaUtils; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; import ru.spcex.platform.enumeration.*; import ru.spcex.platform.utils.enumeration.EnumMessage; @@ -24,6 +33,9 @@ import java.util.Collection; import java.util.Collections; import java.util.Map; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.spy; import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFieldsComparator; class Sdf16ExecutorTest extends AbstractServiceTest { @@ -32,9 +44,14 @@ class Sdf16ExecutorTest extends AbstractServiceTest { private final String acc = "323456789"; @Autowired Sdf16Executor sdf16Executor; + @Captor + ArgumentCaptor producerRecord; private SDf16 sdf16; private StatementRequest statementRequest; + @MockBean + private MockProducer producer; + @PostConstruct void init() { super.init(); @@ -78,11 +95,18 @@ class Sdf16ExecutorTest extends AbstractServiceTest { predictableStatement.setOutSDfId(predictableSdf17.getId()); AccountResult predictableNewResult = getTestAccountResult(currentId.getAndIncrement(), account, company); + MockKafkaUtils.FutureRecordMetadata future = spy(MockKafkaUtils.FutureRecordMetadata.class); + doReturn(future).when(producer).send(producerRecord.capture()); Result result = sdf16Executor.execute(Collections.singletonList(sdf16), statementRequest); + + BaseRequest baseRequest = (BaseRequest) producerRecord.getValue().value(); + NotificationNewRequest notificationNewRequest = (NotificationNewRequest) baseRequest.getRequestPayload(); Statement resultStatement = statementImdg.getSingleObjectByFieldValues(Map.of("account", acc)); SDf17 resultSdf17 = sdf17Imdg.getSingleObjectByFieldValues(Map.of("account", acc)); AccountBalance resultAccountBalance = accountBalanceImdg.getSingleObjectByFieldValues(Map.of("account", acc)); + assertEquals(Consts.NOTIFICATION_NEW, producerRecord.getValue().topic()); + assertEquals(notificationNewRequest.getObjectId(), resultStatement.getId()); RESULT_MATCHER.assertMatch(result, predictableResult); STATEMENT_MATCHER.assertMatch(resultStatement, predictableStatement); SDF_17_MATCHER.assertMatch(resultSdf17, predictableSdf17); diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/StatementServiceServiceTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/StatementServiceServiceTest.java index a9924be94..5e2172dd6 100644 --- a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/StatementServiceServiceTest.java +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/service/StatementServiceServiceTest.java @@ -2,44 +2,57 @@ package ru.spcex.clearing.balance.service; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -import com.hazelcast.core.IMap; -import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.MockConsumer; -import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.clients.producer.MockProducer; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; import org.springframework.beans.factory.annotation.Autowired; -import ru.spcex.clearing.balance.utils.ImapEvent; +import org.springframework.boot.test.mock.mockito.SpyBean; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.sdf.SDf01; import ru.spcex.clearing.platform.messaging.domain.ActionType; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; import ru.spcex.clearing.platform.messaging.service.RequestInfo; -import ru.spcex.clearing.platform.messaging.service.Status; import javax.annotation.PostConstruct; -import java.time.Instant; -import java.util.Collection; -import java.util.Collections; -import java.util.HashMap; -import java.util.Map; +import java.util.concurrent.atomic.AtomicLong; -import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.mockito.Mockito.timeout; +import static org.mockito.Mockito.verify; +import static ru.spcex.clearing.balance.utils.MockKafkaUtils.addRecordToKafka; import static ru.spcex.platform.enumeration.SdfTable.SDF_01; class StatementServiceServiceTest extends AbstractServiceTest { - // private static final MatcherFactory.Matcher SDF_08_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id"); private static final String TOPIC = Consts.STATEMENT_PROCESS; private static final int PARTITION = 1; private static final Long groupId = 111L; - private static final Long limit = 300L; + private final static AtomicLong cuurentOffset = new AtomicLong(1L); + private final Long ID = 11L; @Autowired StatementService statementService; - @Autowired + + @Captor + ArgumentCaptor producerRecord; + private MockConsumer mockConsumer; + private SDf01 sDf01; + private Company company; + + @SpyBean + private MockProducer producer; @PostConstruct void init() { super.init(); + mockConsumer = (MockConsumer) statementService.getConsumer(); + sDf01 = getTestSdf01(ID, groupId); + company = getTestCompany(); } /** @@ -50,43 +63,86 @@ class StatementServiceServiceTest extends AbstractServiceTest { * {@link StatementRequest#table} - SDF_01
*/ @Test - void process() throws InterruptedException { + void processEXPORT_PROCESS() { + sdf01Imdg.delete(sDf01); + companyImdg.delete(company); + + companyImdg.delete(company); StatementRequest statementRequest = new StatementRequest(); statementRequest.setGroupId(groupId); statementRequest.setTable(SDF_01); - BaseRequest baseRequest = new BaseRequest<>(); - baseRequest.setRequestPayload(statementRequest); - baseRequest.setId(currentId.getAndIncrement()); - baseRequest.setActionType(ActionType.NEW); + BaseRequest baseNewRequest = new BaseRequest<>(); + baseNewRequest.setRequestPayload(statementRequest); + baseNewRequest.setId(currentId.getAndIncrement()); + baseNewRequest.setActionType(ActionType.NEW); String jsonBaseNewRequest; ObjectMapper objectMapper = new ObjectMapper(); try { - jsonBaseNewRequest = objectMapper.writeValueAsString(baseRequest); + jsonBaseNewRequest = objectMapper.writeValueAsString(baseNewRequest); } catch (JsonProcessingException e) { throw new RuntimeException(e); } - IMap resultsRequestMap = requestInfoImdg.getMap(); - ImapEvent imapEvent = new ImapEvent(resultsRequestMap); - long timeNow = Instant.now().toEpochMilli(); - //KAFKA - mockConsumer.schedulePollTask(() -> { - mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC, PARTITION))); - mockConsumer.addRecord(new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonBaseNewRequest)); - }); + addRecordToKafka(mockConsumer, TOPIC, PARTITION, cuurentOffset.getAndIncrement(), jsonBaseNewRequest); - HashMap startOffsets = new HashMap<>(); - TopicPartition tp = new TopicPartition(TOPIC, PARTITION); - startOffsets.put(tp, 0L); - mockConsumer.updateBeginningOffsets(startOffsets); + //waiting for kafka producer send message (finale event) + verify(producer, timeout(30_000L).times(1)) + .send(producerRecord.capture()); - //waiting for hazelcast map item updates - imapEvent.waitWhenHappened(); +// Collection resultsRequestInfo = requestInfoImdg.getCollectionObjectsByFieldValues(Map.of("status", Status.Processing)); +// RequestInfo requestInfo = resultsRequestInfo.stream().max((entry1, entry2) -> entry1.getId() > entry2.getId() ? 1 : -1).get(); - Collection resultsRequestInfo = requestInfoImdg.getCollectionObjectsByFieldValues(Map.of("status", Status.Processing)); - RequestInfo requestInfo = resultsRequestInfo.stream().max((entry1, entry2) -> entry1.getId() > entry2.getId() ? 1 : -1).get(); - Long diffRequestInfo = requestInfo.getCreated().toEpochMilli() - timeNow; + assertEquals(Consts.EXPORT_PROCESS, producerRecord.getValue().topic()); + BaseRequest baseRequest = (BaseRequest) producerRecord.getValue().value(); - assertTrue(limit.compareTo(diffRequestInfo) > 0); + RequestInfo resultRequestInfo = requestInfoImdg.getSingleObjectByID(baseRequest.getId()); + + assertNotNull(baseRequest); + assertNotNull(resultRequestInfo); + } + + /** + * {@link StatementService}
+ * Тест проверяет генерацию сущностей {@link RequestInfo}
+ * Входные параметры:
+ * {@link StatementRequest}: new StatementRequest()
+ * {@link StatementRequest#table} - SDF_01
+ */ + @Test + void processACCOUNT_NEW() { + sdf01Imdg.insert(sDf01); + companyImdg.insert(company); + + StatementRequest statementRequest = new StatementRequest(); + statementRequest.setGroupId(groupId); + statementRequest.setTable(SDF_01); + BaseRequest baseNewRequest = new BaseRequest<>(); + baseNewRequest.setRequestPayload(statementRequest); + baseNewRequest.setId(currentId.getAndIncrement()); + baseNewRequest.setActionType(ActionType.NEW); + String jsonBaseNewRequest; + ObjectMapper objectMapper = new ObjectMapper(); + try { + jsonBaseNewRequest = objectMapper.writeValueAsString(baseNewRequest); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + + addRecordToKafka(mockConsumer, TOPIC, PARTITION, cuurentOffset.getAndIncrement(), jsonBaseNewRequest); + + //waiting for kafka producer send message (finale event) + verify(producer, timeout(30_000L).times(1)) + .send(producerRecord.capture()); + +// Collection resultsRequestInfo = requestInfoImdg.getCollectionObjectsByFieldValues(Map.of("status", Status.Processing)); +// RequestInfo requestInfo = resultsRequestInfo.stream().max((entry1, entry2) -> entry1.getId() > entry2.getId() ? 1 : -1).get(); + + assertEquals(Consts.ACCOUNT_NEW, producerRecord.getValue().topic()); + BaseRequest baseRequest = (BaseRequest) producerRecord.getValue().value(); + + RequestInfo resultRequestInfo = requestInfoImdg.getSingleObjectByID(baseRequest.getId()); + + assertNotNull(baseRequest); + assertNotNull(resultRequestInfo); } } \ No newline at end of file diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MockKafkaUtils.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MockKafkaUtils.java new file mode 100644 index 000000000..2ec78cd34 --- /dev/null +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/utils/MockKafkaUtils.java @@ -0,0 +1,54 @@ +package ru.spcex.clearing.balance.utils; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.MockConsumer; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.kafka.common.TopicPartition; + +import java.util.Collections; +import java.util.HashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +public class MockKafkaUtils { + public static void addRecordToKafka(MockConsumer mockConsumer, String topic, int partition, long offset, String jsonValue) { + mockConsumer.schedulePollTask(() -> { + mockConsumer.rebalance(Collections.singletonList(new TopicPartition(topic, partition))); + mockConsumer.addRecord(new ConsumerRecord<>(topic, partition, offset, "key", jsonValue)); + }); + + HashMap startOffsets = new HashMap<>(); + TopicPartition tp = new TopicPartition(topic, partition); + startOffsets.put(tp, 0L); + mockConsumer.updateBeginningOffsets(startOffsets); + } + + public static class FutureRecordMetadata implements Future { + @Override + public boolean cancel(boolean mayInterruptIfRunning) { + return false; + } + + @Override + public boolean isCancelled() { + return false; + } + + @Override + public boolean isDone() { + return false; + } + + @Override + public RecordMetadata get() throws InterruptedException, ExecutionException { + return null; + } + + @Override + public RecordMetadata get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { + return null; + } + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index 8fffd91e8..dd1d2d4ac 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -36,15 +36,15 @@ import java.util.function.Function; * утилитный класс для обработки сообщений из очереди */ public class QueueConsumer implements AutoCloseable { + protected final Map> callbacks; private final Logger log = LoggerFactory.getLogger(getClass()); private final AtomicBoolean closed = new AtomicBoolean(false); private final Consumer consumer; - private Producer producer; private final ExecutorService inputExecutor; private final ExecutorService outputExecutor; - protected final Map> callbacks; private final ObjectMapper json; protected boolean supportStartOffsetTimeWindow; + private Producer producer; public QueueConsumer(Consumer kafkaQueue) { this.consumer = kafkaQueue; @@ -58,6 +58,7 @@ public class QueueConsumer implements AutoCloseable { /** * используем этот конструктор, если хотим класть * в кафку "ответ" - информацию о статусе обработки команд + * * @param kafkaQueue * @param kafkaResponseQueue */ @@ -163,12 +164,13 @@ public class QueueConsumer implements AutoCloseable { /** * Валидация запроса + * * @param Класс проверяемого запроса * @return null если ошибок нет */ public RequestInfoUpdate validate(BaseRequest userRequest, - Function validatorBuilder, - IMessageResolver messageResolver) { + Function validatorBuilder, + IMessageResolver messageResolver) { if (validatorBuilder != null) { R req = userRequest.getRequestPayload(); IValidator validator = validatorBuilder.apply(req); @@ -194,5 +196,8 @@ public class QueueConsumer implements AutoCloseable { outputExecutor.shutdown(); } - + // нужен для тестов поскольку у каждого сервиса свой consumer + public Consumer getConsumer() { + return consumer; + } }