Adding StatementServiceServiceTest.
This commit is contained in:
parent
33cdde3b54
commit
a53232a773
8 changed files with 136 additions and 36 deletions
|
|
@ -5,16 +5,13 @@ import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||||
import org.apache.kafka.clients.producer.MockProducer;
|
import org.apache.kafka.clients.producer.MockProducer;
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
import org.apache.kafka.common.serialization.StringSerializer;
|
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.Bean;
|
||||||
import org.springframework.context.annotation.Configuration;
|
import org.springframework.context.annotation.Configuration;
|
||||||
import org.springframework.context.annotation.Scope;
|
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.JsonSerializer;
|
import ru.spcex.clearing.platform.messaging.serialization.JsonSerializer;
|
||||||
|
|
||||||
@Configuration
|
@Configuration
|
||||||
public class KafkaTestConfig {
|
public class KafkaTestConfig {
|
||||||
|
|
||||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
|
||||||
@Bean
|
@Bean
|
||||||
public MockConsumer<String, Object> createTestConsumer() {
|
public MockConsumer<String, Object> createTestConsumer() {
|
||||||
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ import ru.clearing.classes.statics.data.statement.Statement;
|
||||||
import ru.spcex.clearing.balance.config.*;
|
import ru.spcex.clearing.balance.config.*;
|
||||||
import ru.spcex.clearing.balance.utils.MatcherFactory;
|
import ru.spcex.clearing.balance.utils.MatcherFactory;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||||
import ru.spcex.platform.enumeration.*;
|
import ru.spcex.platform.enumeration.*;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
|
import ru.spcex.platform.imdg.iml.hazelcast.adapter.ImdgHazelcast;
|
||||||
|
|
@ -36,9 +37,11 @@ import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFields
|
||||||
Sdf08Service.class,
|
Sdf08Service.class,
|
||||||
Sdf09Executor.class,
|
Sdf09Executor.class,
|
||||||
Sdf16Executor.class,
|
Sdf16Executor.class,
|
||||||
|
StatementService.class,
|
||||||
LoggingService.class,
|
LoggingService.class,
|
||||||
MessagesConfig.class,
|
MessagesConfig.class,
|
||||||
ValidationConfig.class,
|
ValidationConfig.class,
|
||||||
|
SdfExecutorsConfig.class,
|
||||||
BalanceImdgTestConfig.class,
|
BalanceImdgTestConfig.class,
|
||||||
KafkaSenderConfig.class,
|
KafkaSenderConfig.class,
|
||||||
KafkaTestConfig.class})
|
KafkaTestConfig.class})
|
||||||
|
|
@ -47,7 +50,7 @@ public abstract class AbstractServiceTest {
|
||||||
protected static final MatcherFactory.Matcher<Result> RESULT_MATCHER = usingIgnoringFieldsComparator("account.created", "account.updated", "account.clearingDate", "generationId");
|
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 final MatcherFactory.Matcher<Statement> STATEMENT_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "updated", "id");
|
||||||
|
|
||||||
protected final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy");
|
protected final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy hh:mm");
|
||||||
protected final Long accountIdNew = 10L;
|
protected final Long accountIdNew = 10L;
|
||||||
protected final Long addresseeIdNew = 2L;
|
protected final Long addresseeIdNew = 2L;
|
||||||
protected final String deal = "111111111";
|
protected final String deal = "111111111";
|
||||||
|
|
@ -60,6 +63,7 @@ public abstract class AbstractServiceTest {
|
||||||
protected IMap<Long, CompanySymbols> companySymbolsMap;
|
protected IMap<Long, CompanySymbols> companySymbolsMap;
|
||||||
protected IMap<Long, Account> accountMap;
|
protected IMap<Long, Account> accountMap;
|
||||||
protected ImdgHazelcast<Statement> statementImdg;
|
protected ImdgHazelcast<Statement> statementImdg;
|
||||||
|
protected ImdgHazelcast<RequestInfo> requestInfoImdg;
|
||||||
protected ImdgHazelcast<SDf02> sdf02Imdg;
|
protected ImdgHazelcast<SDf02> sdf02Imdg;
|
||||||
protected ImdgHazelcast<SDf08> sdf08Imdg;
|
protected ImdgHazelcast<SDf08> sdf08Imdg;
|
||||||
protected ImdgHazelcast<SDf10> sdf10Imdg;
|
protected ImdgHazelcast<SDf10> sdf10Imdg;
|
||||||
|
|
@ -74,6 +78,7 @@ public abstract class AbstractServiceTest {
|
||||||
ImdgHazelcast<Company> companyImdg = (ImdgHazelcast<Company>) hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
ImdgHazelcast<Company> companyImdg = (ImdgHazelcast<Company>) hazelcast.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
ImdgHazelcast<CompanySymbols> companySymbolsImdg = (ImdgHazelcast<CompanySymbols>) hazelcast.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.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);
|
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);
|
statementImdg = (ImdgHazelcast<Statement>) hazelcast.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
|
||||||
sdf02Imdg = (ImdgHazelcast<SDf02>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
|
sdf02Imdg = (ImdgHazelcast<SDf02>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
|
||||||
sdf08Imdg = (ImdgHazelcast<SDf08>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class);
|
sdf08Imdg = (ImdgHazelcast<SDf08>) hazelcast.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class);
|
||||||
|
|
|
||||||
|
|
@ -19,7 +19,6 @@ import ru.spcex.platform.utils.number.BigDecimalUtil;
|
||||||
import javax.annotation.PostConstruct;
|
import javax.annotation.PostConstruct;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
import java.time.format.DateTimeFormatter;
|
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
@ -28,7 +27,6 @@ import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFields
|
||||||
|
|
||||||
class Sdf01ExecutorTest extends AbstractServiceTest {
|
class Sdf01ExecutorTest extends AbstractServiceTest {
|
||||||
private static final MatcherFactory.Matcher<SDf02> SDF_02_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
private static final MatcherFactory.Matcher<SDf02> SDF_02_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
||||||
private final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy");
|
|
||||||
private final String acc = "123456789";
|
private final String acc = "123456789";
|
||||||
private final Long ID = 1L;
|
private final Long ID = 1L;
|
||||||
@Autowired
|
@Autowired
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,7 @@
|
||||||
package ru.spcex.clearing.balance.service;
|
package ru.spcex.clearing.balance.service;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||||
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import com.hazelcast.core.IMap;
|
import com.hazelcast.core.IMap;
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||||
|
|
@ -7,9 +9,11 @@ import org.apache.kafka.common.TopicPartition;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf08;
|
import ru.clearing.classes.statics.data.sdf.SDf08;
|
||||||
import ru.clearing.classes.statics.data.statement.Statement;
|
|
||||||
import ru.spcex.clearing.balance.utils.ImapEvent;
|
import ru.spcex.clearing.balance.utils.ImapEvent;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
|
|
||||||
import javax.annotation.PostConstruct;
|
import javax.annotation.PostConstruct;
|
||||||
|
|
@ -17,17 +21,20 @@ import java.time.Instant;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
|
import java.util.Map;
|
||||||
|
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
|
||||||
class Sdf08ServiceTest extends AbstractServiceTest {
|
class Sdf08ServiceTest extends AbstractServiceTest {
|
||||||
// private static final MatcherFactory.Matcher<SDf08> SDF_08_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
// 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 String TOPIC = Task.getAllBalance.topic();
|
||||||
private static final int PARTITION = 1;
|
private static final int PARTITION = 1;
|
||||||
|
private static final Long limit = 300L;
|
||||||
@Autowired
|
|
||||||
private MockConsumer<String, Object> mockConsumer;
|
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
Sdf08Service sdf08Service;
|
Sdf08Service sdf08Service;
|
||||||
|
@Autowired
|
||||||
|
private MockConsumer<String, Object> mockConsumer;
|
||||||
|
|
||||||
@PostConstruct
|
@PostConstruct
|
||||||
void init() {
|
void init() {
|
||||||
|
|
@ -36,26 +43,28 @@ class Sdf08ServiceTest extends AbstractServiceTest {
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* {@link Sdf08Service}<br>
|
* {@link Sdf08Service}<br>
|
||||||
* Тест проверяет генерацию сущностей {@link Result}, {@link Statement}, {@link SDf08}<br>
|
* Тест проверяет генерацию сущностей {@link RequestInfo}<br>
|
||||||
* Входные параметры:<br>
|
* Входные параметры:<br>
|
||||||
* accountId - {@link StatementRequest}: new StatementRequest()<br>
|
* {@link BaseRequest} - new BaseRequest<>()<br>
|
||||||
* addresseeId - {@link Collection<SDf08>}<br>
|
|
||||||
* addresseeId - {@link SDf08}<br>
|
|
||||||
* {@link SDf08#datetime} - текущее время<br>
|
|
||||||
* {@link SDf08#generationId} - id<br>
|
|
||||||
* {@link SDf08#generationTime} - текущее время<br>
|
|
||||||
*/
|
*/
|
||||||
@Test
|
@Test
|
||||||
void newSDf08() throws InterruptedException {
|
void newSDf08() throws InterruptedException {
|
||||||
IMap<Long, SDf08> sdf08Map = sdf08Imdg.getMap();
|
BaseRequest<Object> baseNewRequest = new BaseRequest<>();
|
||||||
SDf08 sDf08 = new SDf08();
|
baseNewRequest.setId(currentId.getAndIncrement());
|
||||||
sDf08.setNumber(idGenerator.nextId().toString());
|
baseNewRequest.setActionType(ActionType.NEW);
|
||||||
Instant now = Instant.now();
|
String jsonBaseNewRequest;
|
||||||
sDf08.setDatetime(String.valueOf(now.toEpochMilli()));
|
ObjectMapper objectMapper = new ObjectMapper();
|
||||||
sDf08.setGenerationTime(now);
|
try {
|
||||||
sDf08.setGenerationId(idGenerator.nextId());
|
jsonBaseNewRequest = objectMapper.writeValueAsString(baseNewRequest);
|
||||||
sdf08Imdg.insert(sDf08);
|
} catch (JsonProcessingException e) {
|
||||||
|
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
|
//KAFKA
|
||||||
mockConsumer.schedulePollTask(() -> {
|
mockConsumer.schedulePollTask(() -> {
|
||||||
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC, PARTITION)));
|
mockConsumer.rebalance(Collections.singletonList(new TopicPartition(TOPIC, PARTITION)));
|
||||||
|
|
@ -68,8 +77,13 @@ class Sdf08ServiceTest extends AbstractServiceTest {
|
||||||
mockConsumer.updateBeginningOffsets(startOffsets);
|
mockConsumer.updateBeginningOffsets(startOffsets);
|
||||||
|
|
||||||
//waiting for hazelcast map item updates
|
//waiting for hazelcast map item updates
|
||||||
ImapEvent imapEvent = new ImapEvent(sdf08Map);
|
|
||||||
imapEvent.waitWhenHappened();
|
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();
|
||||||
|
Long diffRequestInfo = requestInfo.getCreated().toEpochMilli() - timeNow;
|
||||||
|
|
||||||
|
assertEquals(1, sdf08Map.size() - count);
|
||||||
|
assertTrue(limit.compareTo(diffRequestInfo) > 0);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -20,7 +20,6 @@ import javax.annotation.PostConstruct;
|
||||||
import java.math.BigDecimal;
|
import java.math.BigDecimal;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
import java.time.format.DateTimeFormatter;
|
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
@ -29,8 +28,6 @@ import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFields
|
||||||
|
|
||||||
class Sdf09ExecutorTest extends AbstractServiceTest {
|
class Sdf09ExecutorTest extends AbstractServiceTest {
|
||||||
private static final MatcherFactory.Matcher<SDf10> SDF_10_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
private static final MatcherFactory.Matcher<SDf10> SDF_10_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
||||||
|
|
||||||
private final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy");
|
|
||||||
private final Long ID = 2L;
|
private final Long ID = 2L;
|
||||||
private final String acc = "213456789";
|
private final String acc = "213456789";
|
||||||
@Autowired
|
@Autowired
|
||||||
|
|
|
||||||
|
|
@ -20,7 +20,6 @@ import javax.annotation.PostConstruct;
|
||||||
import java.math.BigDecimal;
|
import java.math.BigDecimal;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
import java.time.format.DateTimeFormatter;
|
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
@ -29,8 +28,6 @@ import static ru.spcex.clearing.balance.utils.MatcherFactory.usingIgnoringFields
|
||||||
|
|
||||||
class Sdf16ExecutorTest extends AbstractServiceTest {
|
class Sdf16ExecutorTest extends AbstractServiceTest {
|
||||||
private static final MatcherFactory.Matcher<SDf17> SDF_17_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
private static final MatcherFactory.Matcher<SDf17> SDF_17_MATCHER = usingIgnoringFieldsComparator("created", "comment", "outSDfId", "generationTime", "generationId", "id");
|
||||||
|
|
||||||
private final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy");
|
|
||||||
private final Long ID = 3L;
|
private final Long ID = 3L;
|
||||||
private final String acc = "323456789";
|
private final String acc = "323456789";
|
||||||
@Autowired
|
@Autowired
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,92 @@
|
||||||
|
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.junit.jupiter.api.Test;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
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.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 static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
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;
|
||||||
|
@Autowired
|
||||||
|
StatementService statementService;
|
||||||
|
@Autowired
|
||||||
|
private MockConsumer<String, Object> mockConsumer;
|
||||||
|
|
||||||
|
@PostConstruct
|
||||||
|
void init() {
|
||||||
|
super.init();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link StatementService}<br>
|
||||||
|
* Тест проверяет генерацию сущностей {@link RequestInfo}<br>
|
||||||
|
* Входные параметры:<br>
|
||||||
|
* {@link StatementRequest}: new StatementRequest()<br>
|
||||||
|
* {@link StatementRequest#table} - SDF_01<br>
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void process() throws InterruptedException {
|
||||||
|
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);
|
||||||
|
String jsonBaseNewRequest;
|
||||||
|
ObjectMapper objectMapper = new ObjectMapper();
|
||||||
|
try {
|
||||||
|
jsonBaseNewRequest = objectMapper.writeValueAsString(baseRequest);
|
||||||
|
} 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));
|
||||||
|
});
|
||||||
|
|
||||||
|
HashMap<TopicPartition, Long> startOffsets = new HashMap<>();
|
||||||
|
TopicPartition tp = new TopicPartition(TOPIC, PARTITION);
|
||||||
|
startOffsets.put(tp, 0L);
|
||||||
|
mockConsumer.updateBeginningOffsets(startOffsets);
|
||||||
|
|
||||||
|
//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();
|
||||||
|
Long diffRequestInfo = requestInfo.getCreated().toEpochMilli() - timeNow;
|
||||||
|
|
||||||
|
assertTrue(limit.compareTo(diffRequestInfo) > 0);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -49,7 +49,7 @@ public class ImapEvent<T> {
|
||||||
checkEventHappened.set(secondRan);//если что-то пойдет не так не тормозить основной поток
|
checkEventHappened.set(secondRan);//если что-то пойдет не так не тормозить основной поток
|
||||||
secondRan = true;
|
secondRan = true;
|
||||||
}
|
}
|
||||||
}, 0, 60 * 1000);
|
}, 0, 20 * 1000);
|
||||||
synchronized (checkEventHappened) {
|
synchronized (checkEventHappened) {
|
||||||
while (!checkEventHappened.get()) {
|
while (!checkEventHappened.get()) {
|
||||||
checkEventHappened.wait(100);
|
checkEventHappened.wait(100);
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue