Merge remote-tracking branch 'origin/dev' into dev

This commit is contained in:
ialbert 2023-05-23 21:49:12 +03:00
commit ec619899f6
19 changed files with 199 additions and 333 deletions

View file

@ -32,6 +32,14 @@
<groupId>ru.spcex.clearing</groupId>
<artifactId>classes</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>security-util</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>clearing-validation</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
@ -41,6 +49,11 @@
<artifactId>jackson-databind</artifactId>
</dependency>
<!-- TEST -->
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>test-clearing</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-test</artifactId>
@ -65,14 +78,6 @@
<artifactId>spring-boot-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>security-util</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>clearing-validation</artifactId>
</dependency>
</dependencies>

View file

@ -106,6 +106,7 @@ public class AccountValidationConfig {
Consumer<String> addImdg = (s) -> context.addImdg(s, imdgForValidation.get(s));
addImdg.accept(IMDGDistributedNames.Map_Company);
addImdg.accept(IMDGDistributedNames.Map_Account);
addImdg.accept(IMDGDistributedNames.Map_ServiceStatusDictionary);
addImdg.accept(IMDGDistributedNames.Map_AccountTypeDictionary);
return new ValidatorImpl<>(context,
IdPresentRule.instance("id",

View file

@ -78,13 +78,13 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
@Override
public void afterPropertiesSet() throws Exception {
init();
callback(ClearingAccountNewRequest.class)
.setFunction(this::clearingAccountNew)
.forDestination(Consts.DESTINATION_CLEARING_ACCOUNT_NEW, callbacks::put);
callback(ClearingAccountUpdateRequest.class)
.setFunction(this::clearingAccountUpdate)
.forDestination(Consts.DESTINATION_CLEARING_ACCOUNT_UPDATE, callbacks::put);
init();
}
public RequestInfoUpdate clearingAccountNew(BaseRequest<ClearingAccountNewRequest> userRequest) {

View file

@ -1,70 +0,0 @@
package ru.spcex.clearing.account.config;
import com.hazelcast.config.*;
import com.hazelcast.core.Hazelcast;
import com.hazelcast.core.HazelcastInstance;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import ru.spcex.platform.imdg.iml.hazelcast.util.HazelcastHelper;
import java.util.List;
import java.util.Random;
import java.util.concurrent.atomic.AtomicLong;
@Configuration
public class HazelcastServiceTestConfiguration {
public static final AtomicLong currentID = new AtomicLong(0L);
private HazelcastInstance hazelcastInstance;
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
if (maxPoolSz > 2) {
pool.setKeepAliveSeconds(60);
pool.setAllowCoreThreadTimeOut(true);
}
pool.setCorePoolSize(maxPoolSz);
pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
return pool;
}
@Bean(name = "hazelcastServiceTest")
public HazelcastService hazelcastService(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, HazelcastClientParams params) {
Config cfg = new Config();
cfg.setInstanceName("localhost");
NetworkConfig networkConfig = new NetworkConfig();
JoinConfig joinConfig = new JoinConfig();
joinConfig.setMulticastConfig(new MulticastConfig().setEnabled(false));
joinConfig.setTcpIpConfig(new TcpIpConfig().setEnabled(true).setMembers(List.of("127.0.0.1")));
networkConfig.setJoin(joinConfig);
cfg.setNetworkConfig(networkConfig);
hazelcastInstance = Hazelcast.getOrCreateHazelcastInstance(cfg);
HazelcastHelper.imdgSystem_setStorageState(true, hazelcastInstance);
return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params);
}
@Bean(name = "taskExecutorHazelcastClientInitializer")
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
return createThreadPoolTaskExecutor(1, true);
}
@Bean(name = "taskExecutorIdGeneratorAwaiter")
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
return createThreadPoolTaskExecutor(1, false);
}
@Bean(name = "hazelcastClientParams")
public HazelcastClientParams getHazelcastClientParams() {
HazelcastClientParams params = new HazelcastClientParams();
params.setLogin("dev");
params.setPassword("dev-pass");
params.setClusterMembers("127.0.0.1");
params.setInstanceName("hzTestClient" + new Random().nextInt());
params.setNearCacheConfig(new NearCacheConfig());
return params;
}
}

View file

@ -1,44 +0,0 @@
package ru.spcex.clearing.account.config;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@Configuration
public class KafkaConfigTest {
@Autowired
@Bean(name = "kafkaSenderTest")
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, @Qualifier("hazelcastServiceTest") ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.producer(kafkaProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean(name = "mockConsumerTest")
public MockConsumer<String, Object> createConsumer() {
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
}
}

View file

@ -4,6 +4,7 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@ -21,8 +22,6 @@ import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.platform.dictionary.AccountTypeDictionary;
import ru.clearing.platform.dictionary.ServiceStatusDictionary;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.AccountValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
import ru.spcex.clearing.account.utils.MatcherFactory;
@ -39,6 +38,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
@ -52,9 +54,9 @@ import java.util.UUID;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration.currentID;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.*;
import static ru.spcex.clearing.test.TestUtils.*;
import static ru.spcex.clearing.test.config.ImdgTestConfig.currentID;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -62,8 +64,8 @@ import static ru.spcex.clearing.account.utils.TestUtils.*;
ValidationConfig.class,
AccountValidationConfig.class,
AccountService.class,
HazelcastServiceTestConfiguration.class,
KafkaConfigTest.class})
ImdgTestConfig.class,
KafkaTestConfig.class})
class AccountServiceTest {
public static final MatcherFactory.Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
public static final MatcherFactory.Matcher<RequestInfo> REQUEST_INFO_MATCHER_MATCHER = usingIgnoringFieldsComparator("created");
@ -82,6 +84,9 @@ class AccountServiceTest {
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> producer;
@Autowired
@Qualifier("mockProducer")
protected Producer<String, Object> mockProducer;
private Imdg<Account> accountImdg;
private Imdg<Company> companyImdg;
@ -136,6 +141,8 @@ class AccountServiceTest {
relation.setConsumerId(companyId);
relation.setService(Service.MKR.getKey());
relationImdg.insert(relation);
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
@Test
@ -163,7 +170,7 @@ class AccountServiceTest {
0,
jsonString);
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, mockProducer);
Account resultNew = accountImdg.getCollectionObjectsByFieldValues(Map.of("account", uniqueAccount)).iterator().next();
predictableAccount.setId(resultNew.getId());
@ -190,7 +197,7 @@ class AccountServiceTest {
correspondentAccountUpdateRequest.setId(accountId);
String jsonString = getJsonStringForUPDATE(correspondentAccountUpdateRequest, 0);
String jsonString = getJsonStringForUpdate(correspondentAccountUpdateRequest, 0);
//ACT
addRecordToKafka((MockConsumer) accountService.getConsumer(),
@ -200,7 +207,7 @@ class AccountServiceTest {
jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, mockProducer);
Account resultUpdating = accountImdg.getSingleObjectByID(accountId);
existAccount.setUpdated(resultUpdating.getUpdated());
@ -220,7 +227,7 @@ class AccountServiceTest {
CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
commonDeleteRequest.setId(accountId);
String jsonString = getJsonStringForDELETE(commonDeleteRequest, 0);
String jsonString = getJsonStringForDelete(commonDeleteRequest, 0);
//ACT
addRecordToKafka((MockConsumer) accountService.getConsumer(),
@ -230,7 +237,7 @@ class AccountServiceTest {
jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, mockProducer);
Account resultBlock = accountImdg.getSingleObjectByID(accountId);
existAccount.setStatus(ServiceStatus.Blocked.getKey());
@ -302,6 +309,7 @@ class AccountServiceTest {
//waiting for kafka producer send message (finale event)
verify(producer, timeout(30_000L).times(2))
.send(producerRecord.capture());
//todo переписать валидацию ожидания на новые waitingSendAndCheckRecord / waitingWhenTryAddRecordAndCheckError
ImdgHazelcast<Account> accountImdg = (ImdgHazelcast<Account>) hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Account, Account.class);

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@ -20,8 +21,6 @@ import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.platform.dictionary.CurrencyCodeDictionary;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.AccountValidationConfig;
import ru.spcex.clearing.account.config.validation.BankAccountValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
@ -35,6 +34,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewReq
import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@ -48,8 +50,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.*;
import static ru.spcex.clearing.platform.messaging.service.Status.Error;
import static ru.spcex.clearing.test.TestUtils.*;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -59,8 +61,8 @@ import static ru.spcex.clearing.platform.messaging.service.Status.Error;
BankAccountValidationConfig.class,
AccountValidationConfig.class,
BeanConfiguration.class,
HazelcastServiceTestConfiguration.class,
KafkaConfigTest.class})
ImdgTestConfig.class,
KafkaTestConfig.class})
public class BankAccountServiceTest {
public static final Matcher<BankAccount> BANK_ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
public static final Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
@ -104,6 +106,9 @@ public class BankAccountServiceTest {
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> producer;
@Autowired
@Qualifier("mockProducer")
protected Producer<String, Object> mockProducer;
@PostConstruct
private void init() {
@ -124,6 +129,8 @@ public class BankAccountServiceTest {
currencyCodeDictionary.setCode("RUB");
currencyCodeDictionary.setName("RUB");
currencyCodeDictionaryImdg.insert(currencyCodeDictionary);
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
/**
@ -170,7 +177,7 @@ public class BankAccountServiceTest {
//ACT
addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_NEW, PARTITION, 0, jsonString);
waitingWhenAddedRecordAndCheckIt(ID, producer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
Account accountResult = accountImdg.getSingleObjectBySQL(String.format("account = %s", acc));
BankAccount bankAccountResult = bankAccountImdg.getSingleObjectBySQL(String.format("account = %s or companyId = %s", acc, addresseeIdNew));
@ -343,13 +350,13 @@ public class BankAccountServiceTest {
bankAccountUpdateRequest.setTaxRegistrationReasonCode(predictableUpdateBankAccount.getTaxRegistrationReasonCode());
bankAccountUpdateRequest.setAccount(predictableUpdateBankAccount.getAccount());
String jsonString = getJsonStringForUPDATE(bankAccountUpdateRequest, ID);
String jsonString = getJsonStringForUpdate(bankAccountUpdateRequest, ID);
//ACT
addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_UPDATE, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, producer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
Account accountResult = accountImdg.getSingleObjectByID(predictableAccount.getId());
BankAccount resultUpdating = bankAccountImdg.getSingleObjectByID(predictableUpdateBankAccount.getId());
@ -378,13 +385,13 @@ public class BankAccountServiceTest {
CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
commonDeleteRequest.setId(ID);
String jsonString = getJsonStringForDELETE(commonDeleteRequest, ID);
String jsonString = getJsonStringForDelete(commonDeleteRequest, ID);
//ACT
addRecordToKafka((MockConsumer) bankAccountService.getConsumer(), TOPIC_ACCOUNT_DELETE, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, producer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
BankAccount bankAccount = bankAccountImdg.getSingleObjectByID(ID);
Assertions.assertNull(bankAccount);

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@ -20,8 +21,6 @@ import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.platform.dictionary.AccountTypeDictionary;
import ru.clearing.platform.dictionary.ClearingAccountTypeDictionary;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.AccountValidationConfig;
import ru.spcex.clearing.account.config.validation.ClearingAccountValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
@ -30,6 +29,9 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClearingAccountUpdateRequest;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@ -40,7 +42,7 @@ import java.util.Map;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.*;
import static ru.spcex.clearing.test.TestUtils.*;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -50,8 +52,8 @@ import static ru.spcex.clearing.account.utils.TestUtils.*;
AccountValidationConfig.class,
AccountService.class,
ClearingAccountService.class,
HazelcastServiceTestConfiguration.class,
KafkaConfigTest.class})
ImdgTestConfig.class,
KafkaTestConfig.class})
class ClearingAccountServiceTest {
public static final MatcherFactory.Matcher<ClearingAccount> CLEARING_ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
public static final MatcherFactory.Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator("created", "updated");
@ -73,6 +75,9 @@ class ClearingAccountServiceTest {
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> producer;
@Autowired
@Qualifier("mockProducer")
protected Producer<String, Object> mockProducer;
private Imdg<ClearingAccount> clearingAccountImdg;
private Imdg<Account> accountImdg;
@ -136,6 +141,7 @@ class ClearingAccountServiceTest {
relation.setService(Service.MKR.getKey());
relationImdg.insert(relation);
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
@Test
@ -206,10 +212,10 @@ class ClearingAccountServiceTest {
clearingAccountUpdateRequest.setStatus(0);
clearingAccountUpdateRequest.setDeal(TRADING_CODE);
String jsonString = getJsonStringForUPDATE(clearingAccountUpdateRequest, 0L);
String jsonString = getJsonStringForUpdate(clearingAccountUpdateRequest, 0L);
addRecordToKafka((MockConsumer) clearingAccountService.getConsumer(), Consts.DESTINATION_CLEARING_ACCOUNT_UPDATE, PARTITION, 0, jsonString);
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, mockProducer);
Account accountResult = accountImdg.getSingleObjectByID(accountId);
predictableAccount.setId(accountResult.getId());

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@ -23,19 +24,20 @@ import ru.clearing.classes.statics.data.profile.CompanyInfo;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.*;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.ClientCodeValidationConfig;
import ru.spcex.clearing.account.config.validation.TradingClearingRegistryValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
import ru.spcex.clearing.account.utils.MatcherFactory;
import ru.spcex.clearing.account.utils.TestUtils;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.ClientCodeNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.TestUtils;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.TradingClearingRegistryType;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
@ -44,6 +46,7 @@ import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import javax.annotation.PostConstruct;
import static org.junit.jupiter.api.Assertions.*;
import static ru.spcex.clearing.test.TestUtils.waitingSendAndCheckRecord;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -56,9 +59,8 @@ import static org.junit.jupiter.api.Assertions.*;
ValidationConfig.class,
BeanConfiguration.class,
KafkaConfigTest.class,
HazelcastServiceTestConfiguration.class,})
KafkaTestConfig.class,
ImdgTestConfig.class})
class ClientCodeServiceTest {
private static final int PARTITION = 0;
@ -75,7 +77,7 @@ class ClientCodeServiceTest {
private HazelcastService hazelcastServiceTest;
@Captor
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
@Qualifier("mockProducer")
private MockProducer<String, Object> mockProducer;
private Imdg<ClientCode> clientCodeImdg;
@ -149,8 +151,7 @@ class ClientCodeServiceTest {
depoAcc.setAccountId(depoAccount.getId());
depoAccounts.insert(depoAcc);
TestUtils.FutureRecordMetadata future = Mockito.spy(TestUtils.FutureRecordMetadata.class);
Mockito.doReturn(future).when(mockProducer).send(producerRecord.capture());
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
private <D extends AbstractDictionary> void putToDictionary(String mapName, D object, String code) {
@ -208,7 +209,7 @@ class ClientCodeServiceTest {
TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString);
//ASSERT
TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
@ -246,7 +247,7 @@ class ClientCodeServiceTest {
TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW_UM_COMPANY, PARTITION, 0, jsonString);
//ASSERT
TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
TestUtils.waitingSendAndCheckRecord(ID, mockProducer, producerRecord);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
@ -286,7 +287,7 @@ class ClientCodeServiceTest {
TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_NEW, PARTITION, 0, jsonString);
//ASSERT
TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultNew = clientCodeImdg.getSingleObjectBySQL(String.format("code = '%s'", ccCode));
predictableClientCode.setId(resultNew.getId());
CLIENT_CODE_MATCHER.assertMatch(resultNew, predictableClientCode);
@ -330,12 +331,12 @@ class ClientCodeServiceTest {
predictableClientCode.setStatus("ACTV");
//ACT
String jsonString = TestUtils.getJsonStringForUPDATE(clientCodeUpdateRequest, ID);
String jsonString = TestUtils.getJsonStringForUpdate(clientCodeUpdateRequest, ID);
TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_UPDATE, PARTITION, 0, jsonString);
//ASSERT
TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultUpdating = clientCodeImdg.getSingleObjectByID(ID);
CLIENT_CODE_MATCHER.assertMatch(resultUpdating, predictableClientCode);
@ -369,13 +370,13 @@ class ClientCodeServiceTest {
Assertions.assertNotNull(clientCodeImdg.getSingleObjectByID(ID)); // verify test data
//ACT
String jsonString = TestUtils.getJsonStringForUPDATE(clientCodeDeleteRequest, ID);
String jsonString = TestUtils.getJsonStringForUpdate(clientCodeDeleteRequest, ID);
TestUtils.addRecordToKafka((MockConsumer) clientCodeService.getConsumer(), Consts.DESTINATION_CLIENT_CODE_DELETE, PARTITION, 0, jsonString);
//ASSERT
TestUtils.waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
waitingSendAndCheckRecord(ID, mockProducer);
ClientCode resultUpdate = clientCodeImdg.getSingleObjectByID(ID);
Assertions.assertNull(resultUpdate);
}

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@ -20,8 +21,6 @@ import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.platform.dictionary.AccountTypeDictionary;
import ru.clearing.platform.dictionary.DepoAccountTypeDictionary;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.AccountValidationConfig;
import ru.spcex.clearing.account.config.validation.DepoAccountValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
@ -29,6 +28,9 @@ import ru.spcex.clearing.account.utils.MatcherFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.account.DepoAccountNewRequest;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@ -39,8 +41,8 @@ import java.util.Map;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.addRecordToKafka;
import static ru.spcex.clearing.account.utils.TestUtils.getJsonStringForNew;
import static ru.spcex.clearing.test.TestUtils.addRecordToKafka;
import static ru.spcex.clearing.test.TestUtils.getJsonStringForNew;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -50,8 +52,8 @@ import static ru.spcex.clearing.account.utils.TestUtils.getJsonStringForNew;
AccountValidationConfig.class,
AccountService.class,
DepoAccountService.class,
HazelcastServiceTestConfiguration.class,
KafkaConfigTest.class})
ImdgTestConfig.class,
KafkaTestConfig.class})
class DepoAccountServiceTest {
public static final MatcherFactory.Matcher<DepoAccount> CLEARING_ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
public static final MatcherFactory.Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator("created", "updated");
@ -73,6 +75,9 @@ class DepoAccountServiceTest {
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> producer;
@Autowired
@Qualifier("mockProducer")
protected Producer<String, Object> mockProducer;
private Imdg<DepoAccount> depoAccountImdg;
private Imdg<Account> accountImdg;
@ -136,6 +141,7 @@ class DepoAccountServiceTest {
relation.setService(Service.MKR.getKey());
relationImdg.insert(relation);
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
@Test

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@ -19,8 +20,6 @@ import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.platform.dictionary.AccountTypeDictionary;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.AccountValidationConfig;
import ru.spcex.clearing.account.config.validation.InformationAccountValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
@ -28,6 +27,9 @@ import ru.spcex.clearing.account.utils.MatcherFactory;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.account.InformationAccountNewRequest;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
@ -38,8 +40,8 @@ import java.util.Map;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.addRecordToKafka;
import static ru.spcex.clearing.account.utils.TestUtils.getJsonStringForNew;
import static ru.spcex.clearing.test.TestUtils.addRecordToKafka;
import static ru.spcex.clearing.test.TestUtils.getJsonStringForNew;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -49,8 +51,8 @@ import static ru.spcex.clearing.account.utils.TestUtils.getJsonStringForNew;
AccountValidationConfig.class,
AccountService.class,
InformationAccountService.class,
HazelcastServiceTestConfiguration.class,
KafkaConfigTest.class})
ImdgTestConfig.class,
KafkaTestConfig.class})
class InformationAccountServiceTest {
public static final MatcherFactory.Matcher<InformationAccount> INFORMATION_ACCOUNT_MATCHER = usingIgnoringFieldsComparator();
public static final MatcherFactory.Matcher<Account> ACCOUNT_MATCHER = usingIgnoringFieldsComparator("created", "updated");
@ -71,6 +73,9 @@ class InformationAccountServiceTest {
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> producer;
@Autowired
@Qualifier("mockProducer")
protected Producer<String, Object> mockProducer;
private Imdg<InformationAccount> informationAccountImdg;
private Imdg<Account> accountImdg;
@ -128,6 +133,8 @@ class InformationAccountServiceTest {
accountAnlt.setAccountType(AccountType.Anlt.getKey());
accountAnlt.setCompanyId(1L);
anltAccountId = accountImdg.insert(accountAnlt);
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
@Test

View file

@ -2,6 +2,7 @@ package ru.spcex.clearing.account.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@ -20,8 +21,6 @@ import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.platform.dictionary.ServiceStatusDictionary;
import ru.spcex.clearing.account.config.BeanConfiguration;
import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.account.config.KafkaConfigTest;
import ru.spcex.clearing.account.config.validation.TradingClearingRegistryValidationConfig;
import ru.spcex.clearing.account.config.validation.ValidationConfig;
import ru.spcex.clearing.account.utils.MatcherFactory;
@ -30,6 +29,9 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
import ru.spcex.clearing.test.TestObjectCreator;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.ServiceStatus;
import ru.spcex.platform.enumeration.TradingClearingRegistryPurpose;
import ru.spcex.platform.enumeration.TradingClearingRegistryType;
@ -40,7 +42,7 @@ import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import javax.annotation.PostConstruct;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.account.utils.TestUtils.*;
import static ru.spcex.clearing.test.TestUtils.*;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
@ -48,8 +50,8 @@ import static ru.spcex.clearing.account.utils.TestUtils.*;
ValidationConfig.class,
TradingClearingRegistryValidationConfig.class,
TradingClearingRegistryService.class,
HazelcastServiceTestConfiguration.class,
KafkaConfigTest.class})
ImdgTestConfig.class,
KafkaTestConfig.class})
class TradingClearingRegistryServiceTest {
public static final MatcherFactory.Matcher<TradingClearingRegistry> TRADING_CLEARING_REGISTRY_MATCHER = usingIgnoringFieldsComparator();
private static final int PARTITION = 0;
@ -67,6 +69,9 @@ class TradingClearingRegistryServiceTest {
private ArgumentCaptor<ProducerRecord> producerRecord;
@SpyBean
private MockProducer<String, Object> producer;
@Autowired
@Qualifier("mockProducer")
protected Producer<String, Object> mockProducer;
private Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
private Imdg<Company> companyImdg;
@ -142,6 +147,8 @@ class TradingClearingRegistryServiceTest {
InformationAccount informationAccount = new InformationAccount();
informationAccount.setAccountId(account2Id);
infoAccountId = informationAccountImdg.insert(informationAccount);
new TestObjectCreator(hazelcastServiceTest).createUserAdmin(1000L);
}
@Test
@ -166,7 +173,7 @@ class TradingClearingRegistryServiceTest {
0,
jsonString);
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, mockProducer);
TradingClearingRegistry resultNew = tradingClearingRegistryImdg.getAllValues().iterator().next();
predictableTradingClearingRegistry.setId(resultNew.getId());
@ -202,7 +209,7 @@ class TradingClearingRegistryServiceTest {
0,
jsonString);
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, producer, producerRecord);
TradingClearingRegistry resultNew = tradingClearingRegistryImdg.getAllValues().iterator().next();
predictableTradingClearingRegistry.setId(resultNew.getId());
@ -241,7 +248,7 @@ class TradingClearingRegistryServiceTest {
0,
jsonString);
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, producer, producerRecord);
TradingClearingRegistry resultNew = tradingClearingRegistryImdg.getAllValues().iterator().next();
predictableTradingClearingRegistry.setId(resultNew.getId());
@ -265,7 +272,7 @@ class TradingClearingRegistryServiceTest {
tradingClearingRegistryUpdateRequest.setStatus(ServiceStatus.Active.getKey());
tradingClearingRegistryUpdateRequest.setId(registryId);
String jsonString = getJsonStringForUPDATE(tradingClearingRegistryUpdateRequest, 0);
String jsonString = getJsonStringForUpdate(tradingClearingRegistryUpdateRequest, 0);
//ACT
addRecordToKafka((MockConsumer) tradingClearingRegistryService.getConsumer(),
@ -275,7 +282,7 @@ class TradingClearingRegistryServiceTest {
jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, producer, producerRecord);
TradingClearingRegistry resultUpdating = tradingClearingRegistryImdg.getSingleObjectByID(registryId);
existTradingClearingRegistry.setUpdated(resultUpdating.getUpdated());
@ -293,7 +300,7 @@ class TradingClearingRegistryServiceTest {
CommonDeleteRequest tradingClearingRegistryDeleteRequest = new CommonDeleteRequest();
tradingClearingRegistryDeleteRequest.setId(registryId);
String jsonString = getJsonStringForDELETE(tradingClearingRegistryDeleteRequest, 0);
String jsonString = getJsonStringForDelete(tradingClearingRegistryDeleteRequest, 0);
//ACT
addRecordToKafka((MockConsumer) tradingClearingRegistryService.getConsumer(),
@ -303,7 +310,7 @@ class TradingClearingRegistryServiceTest {
jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(0L, producer, producerRecord);
waitingSendAndCheckRecord(0L, producer, producerRecord);
TradingClearingRegistry resultUpdating = tradingClearingRegistryImdg.getSingleObjectByID(registryId);
existTradingClearingRegistry.setUpdated(resultUpdating.getUpdated());

View file

@ -1,125 +0,0 @@
package ru.spcex.clearing.account.utils;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.TopicPartition;
import org.mockito.ArgumentCaptor;
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.service.RequestInfoUpdate;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.Imdg;
import java.util.Collection;
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;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.platform.messaging.service.Status.Success;
public class TestUtils {
public static final MatcherFactory.Matcher<BaseRequest<Object>> BASE_REQUEST_MATCHER = usingIgnoringFieldsComparator();
private static final ObjectMapper objectMapper = new ObjectMapper();
public static void waitingWhenAddedRecordAndCheckIt(Long id, MockProducer mockProducer, ArgumentCaptor<ProducerRecord> producerRecord) {
BaseRequest<Object> predictableBaseRequest = new BaseRequest<>();
predictableBaseRequest.setId(id);
predictableBaseRequest.setActionType(ActionType.SYSTEM);
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
requestInfoUpdate.setId(id);
requestInfoUpdate.setStatus(Success);
predictableBaseRequest.setRequestPayload(requestInfoUpdate);
//waiting for kafka producer send message (finale event)
verify(mockProducer, timeout(30_000L).times(1))
.send(producerRecord.capture());
BaseRequest<Object> baseRequestResult = (BaseRequest<Object>) producerRecord.getValue().value();
assertEquals(Consts.REQUEST_INFO_UPDATE, producerRecord.getValue().topic());
BASE_REQUEST_MATCHER.assertMatch(baseRequestResult, predictableBaseRequest);
}
public static void addRecordToKafka(MockConsumer mockConsumer, String topic, int partition, long offset, String jsonValue) {
TopicPartition tp = new TopicPartition(topic, partition);
HashMap<TopicPartition, Long> startOffsets = new HashMap<>();
startOffsets.put(tp, 0L);
mockConsumer.updateBeginningOffsets(startOffsets);
mockConsumer.schedulePollTask(() -> {
mockConsumer.rebalance(Collections.singletonList(tp));
mockConsumer.addRecord(new ConsumerRecord<>(topic, partition, offset, "key", jsonValue));
});
}
public static <T> String getJsonStringForNew(T accountRequest, long id) {
return getJsonBaseRequest(accountRequest, id, ActionType.NEW);
}
public static <T> String getJsonStringForUPDATE(T accountRequest, long id) {
return getJsonBaseRequest(accountRequest, id, ActionType.UPDATE);
}
public static <T> String getJsonStringForDELETE(T accountRequest, long id) {
return getJsonBaseRequest(accountRequest, id, ActionType.DELETE);
}
private static <T> String getJsonBaseRequest(T accountRequest, long id, ActionType actionType) {
BaseRequest<T> baseRequest = new BaseRequest<>();
baseRequest.setRequestPayload(accountRequest);
baseRequest.setId(id);
baseRequest.setActionType(actionType);
String jsonBaseRequest;
try {
jsonBaseRequest = objectMapper.writeValueAsString(baseRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
return jsonBaseRequest;
}
public static <T extends SpcexObjectBase> void clearImdg(Imdg<T> imdg) {
Collection<T> values = imdg.getAllValues();
for (T val : values) {
imdg.delete(val);
}
}
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;
}
}
}

View file

@ -36,16 +36,16 @@ public class MoneyExporterService extends AbstractExporterService {
log.debug("Started loading and formation of money file lines");
LocalDate currentDate = LocalDate.now();
List<String> limFileRows = new ArrayList<>();
Collection<Registry> registriesA = registryImdg.getCollectionObjectsByFieldValues(
Map.of("registryDesignation", "A",
"registryInstrumentType", "M",
"registryCode", "F",
"clearingDate", currentDate));
Collection<Registry> registriesD = registryImdg.getCollectionObjectsByFieldValues(
Map.of("registryDesignation", "D",
"registryInstrumentType", "M",
"registryCode", "T",
"clearingDate", currentDate));
Collection<Registry> registriesA = registryImdg.getCollectionObjectsByFieldValues(Map.of(
"registryDesignation", "A",
"registryInstrumentType", "M",
"registryUnit", "F"
));
Collection<Registry> registriesD = registryImdg.getCollectionObjectsByFieldValues(Map.of(
"registryDesignation", "D",
"registryInstrumentType", "M",
"registryUnit", "T"
));
Map<String, List<Registry>> byTcrA = registriesA.stream()
.collect(Collectors.groupingBy(Registry::getTradingClearingRegistry));
Map<String, List<Registry>> byTcrD = registriesD.stream()

View file

@ -35,11 +35,11 @@ public class SecurityExporterService extends AbstractExporterService {
log.debug("Started loading and formation of DEPO file lines");
LocalDate currentDate = LocalDate.now();
List<String> limFileRows = new ArrayList<>();
Collection<Registry> registries = registryImdg.getCollectionObjectsByFieldValues(
Map.of("registryDesignation", "A",
"registryInstrumentType", "S",
"registryCode", "T",
"clearingDate", currentDate));
Collection<Registry> registries = registryImdg.getCollectionObjectsByFieldValues(Map.of(
"registryDesignation", "A",
"registryInstrumentType", "S",
"registryUnit", "T"
));
for (Registry registry : registries) {
if (checkNotBlocked(registry)) {
limFileRows.add(getRow(registry));

View file

@ -0,0 +1,41 @@
package ru.spcex.clearing.test;
import ru.clearing.classes.statics.data.user.User;
import ru.clearing.classes.statics.data.user.UserRoleSession;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.UserRole;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
public class TestObjectCreator {
protected ImdgProvider hazelcastServiceTest;
public TestObjectCreator(ImdgProvider hazelcastServiceTest) {
this.hazelcastServiceTest = hazelcastServiceTest;
}
/**
* Создаёт пользователя администратора
* @param userId
*/
public void createUserAdmin(Long userId) {
Imdg<User> userImdg = hazelcastServiceTest.getImdg(
IMDGDistributedNames.Map_User, User.class
);
User user = new User();
user.setId(userId);
user.setIdentifier("test user");
userImdg.insert(user);
Imdg<UserRoleSession> userRoleSessionImdg = hazelcastServiceTest.getImdg(
IMDGDistributedNames.Map_UserRoleSession, UserRoleSession.class
);
UserRoleSession userRoleSession = new UserRoleSession();
userRoleSession.setId(userId);
userRoleSession.setUserId(user.getId());
userRoleSession.setCompanyId(userId);
userRoleSession.setUserRole(UserRole.Admin.getKey());
userRoleSession.setStatus(WorkflowStatus.Active.getKey());
userRoleSessionImdg.insert(userRoleSession);
}
}

View file

@ -19,6 +19,7 @@ import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
import java.lang.reflect.Field;
import java.util.Collection;
@ -200,4 +201,15 @@ public class TestUtils {
return null;
}
}
public static <T extends SpcexObjectBase> void clearImdg(Imdg<T> imdg) {
if (imdg instanceof ImdgHazelcast) {
((ImdgHazelcast)imdg).clear(); // test API
return;
}
Collection<T> values = imdg.getAllValues();
for (T val : values) {
imdg.delete(val);
}
}
}

View file

@ -100,10 +100,12 @@ public class KafkaTestConfig {
}
@Bean
public KafkaSender kafkaSender(ImdgProvider imdgProvider, Producer<String, Object> mockProducer) {
public KafkaSender kafkaSender(KafkaTemplate<String, Object> kafkaTemplate,
ImdgProvider imdgProvider, Producer<String, Object> mockProducer) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
.producer(mockProducer)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {

View file

@ -75,6 +75,7 @@ public class QueueConsumer implements AutoCloseable {
}
public void init() {
final Exception debugCallerStacktrace = new Exception("init call");
inputExecutor.submit(() -> {
try {
log.debug("Using kafka consumer {} for subscribe on \"{}\"", consumer, callbacks.keySet());
@ -128,7 +129,8 @@ public class QueueConsumer implements AutoCloseable {
} catch (WakeupException e) {
if (!closed.get()) throw e;
} catch (Throwable e) {
log.error(ExceptionUtils.getStackTrace(e));
log.error("QueueConsumer error: {}; main thread stacktrace: {}",
ExceptionUtils.getStackTrace(e), ExceptionUtils.getStackTrace(debugCallerStacktrace));
} finally {
consumer.close();
}