Add Mock kafka producer. Fixing test in balance-service.
This commit is contained in:
parent
2ce4edee96
commit
a915121c3f
14 changed files with 288 additions and 132 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -56,6 +56,16 @@
|
|||
<artifactId>assertj-core</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.mockito</groupId>
|
||||
<artifactId>mockito-core</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
<build>
|
||||
<resources>
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
|
@ -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<String, Object> createTestConsumer() {
|
||||
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Producer<String, Object> createTestProducer() {
|
||||
return new MockProducer<>(true, new StringSerializer(), new JsonSerializer());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<AccountBalance> ACCOUNT_BALANCE_MATCHER = usingIgnoringFieldsComparator("created", "updated", "clearingDate", "id");
|
||||
protected static final MatcherFactory.Matcher<Result> RESULT_MATCHER = usingIgnoringFieldsComparator("account.created", "account.updated", "account.clearingDate", "generationId");
|
||||
protected static final MatcherFactory.Matcher<Statement> STATEMENT_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "updated", "id");
|
||||
protected static AtomicLong currentId = new AtomicLong(0L);
|
||||
protected static IMap<Long, Company> companyMap;
|
||||
protected static IMap<Long, CompanySymbols> companySymbolsMap;
|
||||
protected static IMap<Long, Account> accountMap;
|
||||
protected static ImdgHazelcast<Statement> statementImdg;
|
||||
protected static ImdgHazelcast<RequestInfo> requestInfoImdg;
|
||||
protected static ImdgHazelcast<SDf02> sdf02Imdg;
|
||||
protected static ImdgHazelcast<SDf01> sdf01Imdg;
|
||||
protected static ImdgHazelcast<SDf08> sdf08Imdg;
|
||||
protected static ImdgHazelcast<SDf10> sdf10Imdg;
|
||||
protected static ImdgHazelcast<SDf17> sdf17Imdg;
|
||||
protected static ImdgHazelcast<AccountBalance> accountBalanceImdg;
|
||||
protected static ImdgHazelcast<Company> 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<Long, Company> companyMap;
|
||||
protected IMap<Long, CompanySymbols> companySymbolsMap;
|
||||
protected IMap<Long, Account> accountMap;
|
||||
protected ImdgHazelcast<Statement> statementImdg;
|
||||
protected ImdgHazelcast<RequestInfo> requestInfoImdg;
|
||||
protected ImdgHazelcast<SDf02> sdf02Imdg;
|
||||
protected ImdgHazelcast<SDf08> sdf08Imdg;
|
||||
protected ImdgHazelcast<SDf10> sdf10Imdg;
|
||||
protected ImdgHazelcast<SDf17> sdf17Imdg;
|
||||
protected ImdgHazelcast<AccountBalance> accountBalanceImdg;
|
||||
@Autowired
|
||||
@Qualifier("hazelcastServiceTest")
|
||||
private ImdgProvider hazelcast;
|
||||
|
||||
void init() {
|
||||
hazelcast.waitAvailable();
|
||||
ImdgHazelcast<Company> companyImdg = (ImdgHazelcast<Company>) hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
||||
companyImdg = (ImdgHazelcast<Company>) hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
||||
ImdgHazelcast<CompanySymbols> companySymbolsImdg = (ImdgHazelcast<CompanySymbols>) hazelcast.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
|
||||
ImdgHazelcast<Account> accountImdg = (ImdgHazelcast<Account>) hazelcast.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
||||
requestInfoImdg = (ImdgHazelcast<RequestInfo>) hazelcast.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
|
||||
statementImdg = (ImdgHazelcast<Statement>) hazelcast.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
|
||||
sdf01Imdg = (ImdgHazelcast<SDf01>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
|
||||
sdf02Imdg = (ImdgHazelcast<SDf02>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
|
||||
sdf08Imdg = (ImdgHazelcast<SDf08>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class);
|
||||
sdf10Imdg = (ImdgHazelcast<SDf10>) 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object> producer;
|
||||
|
||||
@PostConstruct
|
||||
void init() {
|
||||
super.init();
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* {@link AccountBalanceService#createAccountBalance(Long addresseeId, Long accountId, BigDecimal amount, String cashMovementCurrencyCode)}<br>
|
||||
* Тест проверяет генерацию сущности {@link AccountResult}<br>
|
||||
|
|
|
|||
|
|
@ -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<SDf02> 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<String, Object> 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);
|
||||
|
|
|
|||
|
|
@ -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<SDf08> 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> producerRecord;
|
||||
private MockConsumer<String, Object> mockConsumer;
|
||||
@SpyBean
|
||||
private MockProducer<String, Object> producer;
|
||||
|
||||
@PostConstruct
|
||||
void init() {
|
||||
super.init();
|
||||
mockConsumer = (MockConsumer<String, Object>) sdf08Service.getConsumer();
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -60,30 +63,21 @@ class Sdf08ServiceTest extends AbstractServiceTest {
|
|||
throw new RuntimeException(e);
|
||||
}
|
||||
|
||||
IMap<Long, SDf08> sdf08Map = sdf08Imdg.getMap();
|
||||
IMap<Long, RequestInfo> 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<TopicPartition, Long> 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<Object> baseRequest = (BaseRequest<Object>) producerRecord.getValue().value();
|
||||
|
||||
Collection<RequestInfo> 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String, Object> producer;
|
||||
|
||||
@PostConstruct
|
||||
void init() {
|
||||
|
|
|
|||
|
|
@ -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> producerRecord;
|
||||
private SDf16 sdf16;
|
||||
private StatementRequest statementRequest;
|
||||
|
||||
@MockBean
|
||||
private MockProducer<String, Object> 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<Object> baseRequest = (BaseRequest<Object>) 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);
|
||||
|
|
|
|||
|
|
@ -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<SDf08> 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> producerRecord;
|
||||
|
||||
private MockConsumer<String, Object> mockConsumer;
|
||||
private SDf01 sDf01;
|
||||
private Company company;
|
||||
|
||||
@SpyBean
|
||||
private MockProducer<String, Object> producer;
|
||||
|
||||
@PostConstruct
|
||||
void init() {
|
||||
super.init();
|
||||
mockConsumer = (MockConsumer<String, Object>) statementService.getConsumer();
|
||||
sDf01 = getTestSdf01(ID, groupId);
|
||||
company = getTestCompany();
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -50,43 +63,86 @@ class StatementServiceServiceTest extends AbstractServiceTest {
|
|||
* {@link StatementRequest#table} - SDF_01<br>
|
||||
*/
|
||||
@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<StatementRequest> baseRequest = new BaseRequest<>();
|
||||
baseRequest.setRequestPayload(statementRequest);
|
||||
baseRequest.setId(currentId.getAndIncrement());
|
||||
baseRequest.setActionType(ActionType.NEW);
|
||||
BaseRequest<StatementRequest> 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<Long, RequestInfo> 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<TopicPartition, Long> 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<RequestInfo> resultsRequestInfo = requestInfoImdg.getCollectionObjectsByFieldValues(Map.of("status", Status.Processing));
|
||||
// RequestInfo requestInfo = resultsRequestInfo.stream().max((entry1, entry2) -> entry1.getId() > entry2.getId() ? 1 : -1).get();
|
||||
|
||||
Collection<RequestInfo> 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<Object> baseRequest = (BaseRequest<Object>) producerRecord.getValue().value();
|
||||
|
||||
assertTrue(limit.compareTo(diffRequestInfo) > 0);
|
||||
RequestInfo resultRequestInfo = requestInfoImdg.getSingleObjectByID(baseRequest.getId());
|
||||
|
||||
assertNotNull(baseRequest);
|
||||
assertNotNull(resultRequestInfo);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link StatementService}<br>
|
||||
* Тест проверяет генерацию сущностей {@link RequestInfo}<br>
|
||||
* Входные параметры:<br>
|
||||
* {@link StatementRequest}: new StatementRequest()<br>
|
||||
* {@link StatementRequest#table} - SDF_01<br>
|
||||
*/
|
||||
@Test
|
||||
void processACCOUNT_NEW() {
|
||||
sdf01Imdg.insert(sDf01);
|
||||
companyImdg.insert(company);
|
||||
|
||||
StatementRequest statementRequest = new StatementRequest();
|
||||
statementRequest.setGroupId(groupId);
|
||||
statementRequest.setTable(SDF_01);
|
||||
BaseRequest<StatementRequest> 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<RequestInfo> 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<Object> baseRequest = (BaseRequest<Object>) producerRecord.getValue().value();
|
||||
|
||||
RequestInfo resultRequestInfo = requestInfoImdg.getSingleObjectByID(baseRequest.getId());
|
||||
|
||||
assertNotNull(baseRequest);
|
||||
assertNotNull(resultRequestInfo);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<TopicPartition, Long> startOffsets = new HashMap<>();
|
||||
TopicPartition tp = new TopicPartition(topic, partition);
|
||||
startOffsets.put(tp, 0L);
|
||||
mockConsumer.updateBeginningOffsets(startOffsets);
|
||||
}
|
||||
|
||||
public static class FutureRecordMetadata implements Future<RecordMetadata> {
|
||||
@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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -36,15 +36,15 @@ import java.util.function.Function;
|
|||
* утилитный класс для обработки сообщений из очереди
|
||||
*/
|
||||
public class QueueConsumer implements AutoCloseable {
|
||||
protected final Map<String, ConsumerSpecificClass<?>> callbacks;
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final AtomicBoolean closed = new AtomicBoolean(false);
|
||||
private final Consumer<String, Object> consumer;
|
||||
private Producer<String, Object> producer;
|
||||
private final ExecutorService inputExecutor;
|
||||
private final ExecutorService outputExecutor;
|
||||
protected final Map<String, ConsumerSpecificClass<?>> callbacks;
|
||||
private final ObjectMapper json;
|
||||
protected boolean supportStartOffsetTimeWindow;
|
||||
private Producer<String, Object> producer;
|
||||
|
||||
public QueueConsumer(Consumer<String, Object> 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 <R> Класс проверяемого запроса
|
||||
* @return null если ошибок нет
|
||||
*/
|
||||
public <R> RequestInfoUpdate validate(BaseRequest<R> userRequest,
|
||||
Function<R, IValidator> validatorBuilder,
|
||||
IMessageResolver messageResolver) {
|
||||
Function<R, IValidator> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue