Добавил FixedIncomeSecurityService.java, исправил ошибки, добавил тесты.
This commit is contained in:
psemenkov 2023-04-04 18:13:52 +03:00
parent 30485d7104
commit e2ce49e3e5
24 changed files with 1349 additions and 424 deletions

View file

@ -24,6 +24,10 @@
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-enum</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.clearing</groupId>
<artifactId>clearing-utils</artifactId>
</dependency>
<dependency>
<groupId>ru.spcex.platform</groupId>
<artifactId>platform-imdg-api-hazelcast-impl</artifactId>

View file

@ -1,9 +1,11 @@
package ru.spcex.clearing.securities.component;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.FixedIncomeSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.platform.enumeration.Market;
import ru.spcex.platform.enumeration.Status;
@ -23,6 +25,17 @@ public class ListingBuilder {
listing.setWorkflowStatus(Status.Active.getKey());
return listing;
}
public Listing byUpdateMms(MoneyMarketSecurity mms, Listing listing) {
listing.setSecurityId(mms.getId());
listing.setSymbolCode(mms.getSecuritySymbol());
listing.setSymbolName(mms.getFullName());
listing.setTradingCurrency(mms.getNominalCurrency());
listing.setWorkflowStatus(mms.getWorkflowStatus());
listing.setLotSize(mms.getLotSize());
listing.setUpdated(mms.getUpdated());
return listing;
}
public Listing byNewEquity(EquitySecurity equitySecurity, EquitySecurityNewRequest request) {
Listing listing = new Listing();
listing.setCreated(equitySecurity.getCreated());
@ -35,4 +48,27 @@ public class ListingBuilder {
listing.setWorkflowStatus(Status.Active.getKey());
return listing;
}
public Listing byUpdateEquity(EquitySecurity equitySecurity, Listing listing) {
listing.setSecurityId(equitySecurity.getId());
listing.setSymbolCode(equitySecurity.getSecuritySymbol());
listing.setSymbolName(equitySecurity.getFullName());
listing.setWorkflowStatus(equitySecurity.getWorkflowStatus());
listing.setLotSize(equitySecurity.getLotSize());
listing.setUpdated(equitySecurity.getUpdated());
return listing;
}
public Listing byNewFixedIncome(FixedIncomeSecurity instant, FixedIncomeSecurityNewRequest request) {
Listing listing = new Listing();
listing.setCreated(instant.getCreated());
listing.setSecurityId(instant.getId());
listing.setLotSize(BigDecimal.valueOf(request.getLotSize()));
listing.setMarket(Market.mkrs.getKey());
listing.setSymbolCode(instant.getSecuritySymbol());
listing.setSymbolName(instant.getFullName());
listing.setTradingCurrency(instant.getNominalCurrency());
listing.setWorkflowStatus(Status.Active.getKey());
return listing;
}
}

View file

@ -3,8 +3,10 @@ package ru.spcex.clearing.securities.config;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.springframework.beans.factory.annotation.Autowired;
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.platform.messaging.config.KafkaConsumerFactory;
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
import ru.spcex.clearing.securities.config.element.SecuritiesServiceSettings;
@ -13,6 +15,7 @@ import ru.spcex.clearing.securities.config.element.SecuritiesServiceSettings;
public class KafkaConfig {
@Autowired
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Bean
public Consumer<String, Object> createConsumer(SecuritiesServiceSettings settings) {
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());

View file

@ -0,0 +1,17 @@
package ru.spcex.clearing.securities.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
@Configuration
public class ValidationConfig {
@Bean("validationHelper")
public ValidationHelper validationHelper(IMessageResolver messageResolver) {
return new ValidationHelper(messageResolver);
}
}

View file

@ -17,19 +17,17 @@ import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurit
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.clearing.securities.component.ListingBuilder;
import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.ImdgTransaction;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.Optional;
import java.util.Map;
@Service
public class EquitySecurityService extends QueueConsumer implements InitializingBean {
@ -39,53 +37,46 @@ public class EquitySecurityService extends QueueConsumer implements Initializing
private final ImdgProvider imdgProvider;
private final ImdgId idGenerator;
private final ValidationProvider validation;
private final IMessageResolver messageResolver;
private final ValidationHelper validationHelper;
private final ListingBuilder listingBuilder = new ListingBuilder();
@Autowired
public EquitySecurityService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationProvider validation,
IMessageResolver messageResolver) {
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationProvider validation,
ValidationHelper validationHelper) {
super(kafkaQueue, kafkaProducer);
this.validation = validation;
this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
this.imdgProvider = imdgProvider;
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.messageResolver = messageResolver;
this.validationHelper = validationHelper;
}
@Override
public void afterPropertiesSet() {
callback(EquitySecurityNewRequest.class)
.setFunction(this::newEquity)
.forDestination(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, callbacks::put);
.forDestination(Consts.DESTINATION_EQUITY_SECURITY_NEW, callbacks::put);
callback(EquitySecurityUpdateRequest.class)
.setFunction(this::updateEquity)
.forDestination(Consts.DESTINATION_MONEY_MARKET_SECURITY_UPDATE, callbacks::put);
.forDestination(Consts.DESTINATION_EQUITY_SECURITY_UPDATE, callbacks::put);
callback(CommonDeleteRequest.class)
.setFunction(this::deleteEquity)
.forDestination(Consts.DESTINATION_MONEY_MARKET_SECURITY_DELETE, callbacks::put);
.forDestination(Consts.DESTINATION_EQUITY_SECURITY_DELETE, callbacks::put);
imdgProvider.waitAvailable();
init();
}
private RequestInfoUpdate newEquity(BaseRequest<EquitySecurityNewRequest> userRequest){
private RequestInfoUpdate newEquity(BaseRequest<EquitySecurityNewRequest> userRequest) {
ImdgTransaction transaction = imdgProvider.newTransaction();
EquitySecurityNewRequest req = userRequest.getRequestPayload();
Optional<EnumMessage> validationError = validation.equityNewValidator()
.apply(req)
.tillFirstError();
if (validationError.isPresent()) {
String errorMsg = messageResolver.resolve(validationError.get());
log.error("cannot process EquitySecurityNewRequest id={}: {}", userRequest.getId(), errorMsg);
return new RequestInfoUpdate()
.setId(userRequest.getId())
.setStatus(Status.Error)
.setMessage(errorMsg);
}
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.equityNewValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
log.debug("EquitySecurityNewRequest received");
EquitySecurity equity = new EquitySecurity();
equity.setId(idGenerator.nextId());
@ -96,6 +87,7 @@ public class EquitySecurityService extends QueueConsumer implements Initializing
equity.setShareType(req.getShareType());
equity.setLotSize(req.getLotSize() != null ? BigDecimal.valueOf(req.getLotSize()) : null);
equity.setIssuerId(req.getIssuerId());
equity.setShortNameEng(req.getShortNameEng());
equity.setFullNameEng(req.getFullNameEng());
equity.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
equity.setInstrumentType(req.getInstrumentType());
@ -121,11 +113,60 @@ public class EquitySecurityService extends QueueConsumer implements Initializing
return null; // default success
}
private RequestInfoUpdate updateEquity(BaseRequest<EquitySecurityUpdateRequest> userRequest){
private RequestInfoUpdate updateEquity(BaseRequest<EquitySecurityUpdateRequest> userRequest) {
EquitySecurityUpdateRequest req = userRequest.getRequestPayload();
log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId());
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.equityUpdateValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
EquitySecurity equity = equitySecurityImdg.getSingleObjectByID(req.getId());
Instant updateTime = Instant.now();
equity.setUpdated(updateTime);
equity.setSecuritySymbol(req.getSecuritySymbol());
equity.setShortName(req.getShortName());
equity.setFullName(req.getFullName());
equity.setIsin(req.getIsin());
equity.setShareType(req.getShareType());
equity.setLotSize(req.getLotSize() != null ? BigDecimal.valueOf(req.getLotSize()) : null);
equity.setIssuerId(req.getIssuerId());
equity.setShortNameEng(req.getShortNameEng());
equity.setFullNameEng(req.getFullNameEng());
equity.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
equity.setInstrumentType(req.getInstrumentType());
equitySecurityImdg.update(equity);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equity.getId()));
if (listing == null) {
log.error("MoneyMarketSecurityUpdateRequest id {} couldn't find listing with securityId {}", req.getId(), equity.getId());
return null;
}
listing = listingBuilder.byUpdateEquity(equity, listing);
listingImdg.update(listing);
return null;
}
private RequestInfoUpdate deleteEquity(BaseRequest<CommonDeleteRequest> userRequest) {
CommonDeleteRequest req = userRequest.getRequestPayload();
log.debug("CommonDeleteRequest received id = {}", req.getId());
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.equityDeleteValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
Instant updateTime = Instant.now();
EquitySecurity equity = equitySecurityImdg.getSingleObjectByID(req.getId());
equity.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
equity.setUpdated(updateTime);
equitySecurityImdg.update(equity);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equity.getId()));
if (listing == null) {
log.error("CommonDeleteRequest id {} couldn't find listing with securityId {}", req.getId(), equity.getId());
return null;
}
listing.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
listing.setUpdated(updateTime);
listingImdg.update(listing);
return null;
}
}

View file

@ -0,0 +1,189 @@
package ru.spcex.clearing.securities.service.cud;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.clearing.classes.statics.data.misc.Listing;
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.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.FixedIncomeSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.FixedIncomeSecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.securities.component.ListingBuilder;
import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.ImdgTransaction;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.Map;
@Service
public class FixedIncomeSecurityService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<FixedIncomeSecurity> fixedIncomeSecurityImdg;
private final Imdg<Listing> listingImdg;
private final ImdgProvider imdgProvider;
private final ImdgId idGenerator;
private final ValidationProvider validation;
private final ValidationHelper validationHelper;
private final ListingBuilder listingBuilder = new ListingBuilder();
@Autowired
public FixedIncomeSecurityService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationProvider validation,
ValidationHelper validationHelper) {
super(kafkaQueue, kafkaProducer);
this.validation = validation;
this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
this.imdgProvider = imdgProvider;
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.validationHelper = validationHelper;
}
@Override
public void afterPropertiesSet() {
callback(FixedIncomeSecurityNewRequest.class)
.setFunction(this::newFixedIncome)
.forDestination(Consts.DESTINATION_FIXED_INCOME_SECURITY_NEW, callbacks::put);
callback(FixedIncomeSecurityUpdateRequest.class)
.setFunction(this::updateFixedIncome)
.forDestination(Consts.DESTINATION_FIXED_INCOME_SECURITY_UPDATE, callbacks::put);
callback(CommonDeleteRequest.class)
.setFunction(this::deleteFixedIncome)
.forDestination(Consts.DESTINATION_FIXED_INCOME_SECURITY_DELETE, callbacks::put);
imdgProvider.waitAvailable();
init();
}
private RequestInfoUpdate newFixedIncome(BaseRequest<FixedIncomeSecurityNewRequest> userRequest) {
ImdgTransaction transaction = imdgProvider.newTransaction();
FixedIncomeSecurityNewRequest req = userRequest.getRequestPayload();
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.fixedIncomeNewValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
log.debug("EquitySecurityNewRequest received");
FixedIncomeSecurity fixedIncome = new FixedIncomeSecurity();
fixedIncome.setId(idGenerator.nextId());
fixedIncome.setSecuritySymbol(req.getSecuritySymbol());
fixedIncome.setShortName(req.getShortName());
fixedIncome.setFullName(req.getFullName());
fixedIncome.setIsin(req.getIsin());
fixedIncome.setBondType(req.getBondType());
fixedIncome.setLotSize(req.getLotSize() != null ? BigDecimal.valueOf(req.getLotSize()) : null);
fixedIncome.setNominalValue(req.getNominalValue() != null ? BigDecimal.valueOf(req.getNominalValue()) : null);
fixedIncome.setNominalCurrency(req.getNominalCurrency());
fixedIncome.setMaturityDate(req.getMaturityDate());
fixedIncome.setCoupon(req.getCoupon() != null ? BigDecimal.valueOf(req.getCoupon()) : null);
fixedIncome.setCouponFrequency(req.getCouponFrequency());
fixedIncome.setIssuerId(req.getIssuerId());
fixedIncome.setShortNameEng(req.getShortNameEng());
fixedIncome.setFullNameEng(req.getFullNameEng());
fixedIncome.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
fixedIncome.setInstrumentType(req.getInstrumentType());
fixedIncome.setSecurityId(fixedIncome.getId());
fixedIncome.setCreated(Instant.now());
try {
transaction.beginTransaction();
Imdg<FixedIncomeSecurity> fixedIncomeImdg = transaction.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
Imdg<Listing> listingMap = transaction.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
//одним скопом выполняем реквест
fixedIncomeImdg.insert(fixedIncome);
Listing listing = listingBuilder.byNewFixedIncome(fixedIncome, req);
listingMap.insert(listing);
//и сохраняем обновленный
// reqInfoMap.insert(reqInfo);
transaction.commitTransaction();
} catch (Throwable e) {
transaction.rollbackTransaction();
// RequestInfo.update(reqInfo, Status.Error, ExceptionUtils.getStackTrace(e));
// requestInfoImdg.insert(reqInfo);
}
log.debug("successfully processed, request id {}, new object id {}", userRequest.getId(), fixedIncome.getId());
return null; // default success
}
private RequestInfoUpdate updateFixedIncome(BaseRequest<FixedIncomeSecurityUpdateRequest> userRequest) {
FixedIncomeSecurityUpdateRequest req = userRequest.getRequestPayload();
log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId());
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.fixedIncomeUpdateValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
FixedIncomeSecurity fixedIncome = fixedIncomeSecurityImdg.getSingleObjectByID(req.getId());
Instant updateTime = Instant.now();
fixedIncome.setUpdated(updateTime);
fixedIncome.setSecuritySymbol(req.getSecuritySymbol());
fixedIncome.setShortName(req.getShortName());
fixedIncome.setFullName(req.getFullName());
fixedIncome.setIsin(req.getIsin());
fixedIncome.setBondType(req.getBondType());
fixedIncome.setLotSize(req.getLotSize() != null ? BigDecimal.valueOf(req.getLotSize()) : null);
fixedIncome.setNominalValue(req.getNominalValue() != null ? BigDecimal.valueOf(req.getNominalValue()) : null);
fixedIncome.setNominalCurrency(req.getNominalCurrency());
fixedIncome.setMaturityDate(req.getMaturityDate());
fixedIncome.setCoupon(req.getCoupon() != null ? BigDecimal.valueOf(req.getCoupon()) : null);
fixedIncome.setCouponFrequency(req.getCouponFrequency());
fixedIncome.setIssuerId(req.getIssuerId());
fixedIncome.setShortNameEng(req.getShortNameEng());
fixedIncome.setFullNameEng(req.getFullNameEng());
fixedIncome.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
fixedIncome.setInstrumentType(req.getInstrumentType());
fixedIncome.setSecurityId(fixedIncome.getId());
fixedIncomeSecurityImdg.update(fixedIncome);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", fixedIncome.getId()));
if (listing == null) {
log.error("MoneyMarketSecurityUpdateRequest id {} couldn't find listing with securityId {}", req.getId(), fixedIncome.getId());
return null;
}
listing.setSecurityId(fixedIncome.getId());
listing.setSymbolCode(fixedIncome.getSecuritySymbol());
listing.setSymbolName(fixedIncome.getFullName());
listing.setTradingCurrency(fixedIncome.getNominalCurrency());
listing.setWorkflowStatus(fixedIncome.getWorkflowStatus());
listing.setLotSize(fixedIncome.getLotSize());
listing.setUpdated(updateTime);
listingImdg.update(listing);
return null;
}
private RequestInfoUpdate deleteFixedIncome(BaseRequest<CommonDeleteRequest> userRequest) {
CommonDeleteRequest req = userRequest.getRequestPayload();
log.debug("CommonDeleteRequest received id = {}", req.getId());
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.equityDeleteValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
Instant updateTime = Instant.now();
FixedIncomeSecurity equity = fixedIncomeSecurityImdg.getSingleObjectByID(req.getId());
equity.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
equity.setUpdated(updateTime);
fixedIncomeSecurityImdg.update(equity);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equity.getId()));
if (listing == null) {
log.error("CommonDeleteRequest id {} couldn't find listing with securityId {}", req.getId(), equity.getId());
return null;
}
listing.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
listing.setUpdated(updateTime);
listingImdg.update(listing);
return null;
}
}

View file

@ -21,6 +21,7 @@ import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.clearing.securities.component.ListingBuilder;
import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
@ -43,6 +44,7 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
private final ImdgProvider imdgProvider;
private final ImdgId idGenerator;
private final ValidationProvider validation;
private final ValidationHelper validationHelper;
private final IMessageResolver messageResolver;
private final ListingBuilder listingBuilder = new ListingBuilder();
@ -51,8 +53,10 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationProvider validation,
ValidationHelper validationHelper,
IMessageResolver messageResolver) {
super(kafkaQueue, kafkaProducer);
this.validationHelper = validationHelper;
this.validation = validation;
this.moneyMarketSecurityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
@ -94,17 +98,10 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
private RequestInfoUpdate newMoneyMarket(BaseRequest<MoneyMarketSecurityNewRequest> userRequest) {
ImdgTransaction transaction = imdgProvider.newTransaction();
MoneyMarketSecurityNewRequest req = userRequest.getRequestPayload();
Optional<EnumMessage> validationError = validation.mmsNewValidator()
.apply(req)
.tillFirstError();
if (validationError.isPresent()) {
String errorMsg = messageResolver.resolve(validationError.get());
log.error("cannot process MoneyMarketSecurityNewRequest id={}: {}", userRequest.getId(), errorMsg);
return new RequestInfoUpdate()
.setId(userRequest.getId())
.setStatus(Status.Error)
.setMessage(errorMsg);
}
RequestInfoUpdate requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, validation.mmsNewValidator());
if (requestInfoUpdate != null) return requestInfoUpdate;
log.debug("MoneyMarketSecurityNewRequest received");
MoneyMarketSecurity mms = new MoneyMarketSecurity();
mms.setId(idGenerator.nextId());
@ -122,7 +119,7 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
mms.setTermType(req.getTermType());
mms.setIsin(req.getIsin());
mms.setSecurityId(mms.getId());
mms.setDescription(mms.getDescription());
mms.setDescription(req.getDescription());
mms.setCreated(Instant.now());
mms.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
try {
@ -176,7 +173,6 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
mms.setIsin(req.getIsin());
mms.setSecurityId(mms.getId());
mms.setDescription(mms.getDescription());
mms.setCreated(Instant.now());
moneyMarketSecurityMap.update(mms);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", mms.getId()));
@ -184,8 +180,7 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
log.error("MoneyMarketSecurityUpdateRequest id {} couldn't find listing with securityId {}", req.getId(), mms.getId());
return null;
}
listing.setLotSize(mms.getLotSize());
listing.setUpdated(updateTime);
listing = listingBuilder.byUpdateMms(mms, listing);
listingImdg.update(listing);
return null;
}
@ -203,9 +198,10 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
.setStatus(Status.Error)
.setMessage(errorMsg);
}
Instant updateTime = Instant.now();
MoneyMarketSecurity mms = validator.getStored(Stored.PresentById);
mms.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
mms.setUpdated(Instant.now());
mms.setUpdated(updateTime);
moneyMarketSecurityMap.update(mms);
Listing listing = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", mms.getId()));
if (listing == null) {
@ -213,7 +209,7 @@ public class MoneyMarketSecurityService extends QueueConsumer implements Initial
return null;
}
listing.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
listing.setUpdated(Instant.now());
listing.setUpdated(updateTime);
listingImdg.update(listing);
return null;
}

View file

@ -3,10 +3,7 @@ package ru.spcex.clearing.securities.validation;
import org.springframework.stereotype.Component;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.*;
import ru.spcex.clearing.securities.errors.SecuritiesError;
import ru.spcex.clearing.securities.validation.rule.*;
import ru.spcex.platform.classes.base.SpcexObjectBase;
@ -23,11 +20,13 @@ import java.util.function.Function;
public class ValidationProvider {
private final Imdg<SpcexObjectBase> mmsMap;
private final Imdg<SpcexObjectBase> equityMap;
private final Imdg<SpcexObjectBase> fixedIncomeMap;
public ValidationProvider(ImdgProvider imdgProvider) {
this.mmsMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, SpcexObjectBase.class);
this.equityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, SpcexObjectBase.class);
this.fixedIncomeMap = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, SpcexObjectBase.class);
}
public Function<MoneyMarketSecurityNewRequest, IValidator> mmsNewValidator() {
@ -38,7 +37,7 @@ public class ValidationProvider {
return new ValidatorImpl<>(context,
EndDtAfterStartDt.instance(MoneyMarketSecurityNewRequest.class),
new NotPresentBySecuritySymbolAndWorkflowStatusActv(IMDGDistributedNames.Map_MoneyMarketSecurity, SecuritiesError.InstrumentAlreadyExists, MoneyMarketSecurityNewRequest.class),
new IsValidInstrumentType(SecuritiesError.RequiredFieldIsEmpty, "RATE"),
new IsValidInstrumentTypeFromRequest(SecuritiesError.RequiredFieldIsEmpty, "RATE"),
MmsNewValidationRule.RequiredFieldIsNotEmpty
);
};
@ -51,8 +50,8 @@ public class ValidationProvider {
context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, mmsMap);
return new ValidatorImpl<>(context,
EndDtAfterStartDt.instance(MoneyMarketSecurityUpdateRequest.class),
new IsValidInstrumentType(SecuritiesError.InstrumentNotFound, "RATE"),
new PresentById(IMDGDistributedNames.Map_MoneyMarketSecurity, SecuritiesError.InstrumentNotFound, true),
new IsValidInstrumentTypeById(SecuritiesError.InstrumentNotFound, "RATE"),
new WorkflowStatusActv()
);
};
@ -64,8 +63,8 @@ public class ValidationProvider {
context.setValidatedObject(mmsRequest);
context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, mmsMap);
return new ValidatorImpl<>(context,
new IsValidInstrumentType(SecuritiesError.InstrumentNotFound, "RATE"),
new PresentById(IMDGDistributedNames.Map_MoneyMarketSecurity, SecuritiesError.InstrumentNotFound, true),
new IsValidInstrumentTypeById(SecuritiesError.InstrumentNotFound, "RATE"),
new WorkflowStatusActv()
);
};
@ -78,7 +77,7 @@ public class ValidationProvider {
context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equityMap);
return new ValidatorImpl<>(context,
new NotPresentBySecuritySymbolAndWorkflowStatusActv(IMDGDistributedNames.Map_EquitySecurity, SecuritiesError.InstrumentAlreadyExists, EquitySecurityNewRequest.class),
new IsValidInstrumentType(SecuritiesError.RequiredFieldIsEmpty, "EQTY"),
new IsValidInstrumentTypeFromRequest(SecuritiesError.RequiredFieldIsEmpty, "EQTY"),
EquityNewValidationRule.RequiredFieldIsNotEmpty
);
};
@ -90,8 +89,8 @@ public class ValidationProvider {
context.setValidatedObject(equityRequest);
context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equityMap);
return new ValidatorImpl<>(context,
new IsValidInstrumentType(SecuritiesError.InstrumentNotFound, "EQTY"),
new PresentById(IMDGDistributedNames.Map_EquitySecurity, SecuritiesError.InstrumentNotFound, true),
new IsValidInstrumentTypeById(SecuritiesError.InstrumentNotFound, "EQTY"),
new WorkflowStatusActv()
);
};
@ -103,8 +102,47 @@ public class ValidationProvider {
context.setValidatedObject(equityRequest);
context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equityMap);
return new ValidatorImpl<>(context,
new IsValidInstrumentType(SecuritiesError.InstrumentNotFound, "EQTY"),
new PresentById(IMDGDistributedNames.Map_EquitySecurity, SecuritiesError.InstrumentNotFound, true),
new IsValidInstrumentTypeById(SecuritiesError.InstrumentNotFound, "EQTY"),
new WorkflowStatusActv()
);
};
}
public Function<FixedIncomeSecurityNewRequest, IValidator> fixedIncomeNewValidator() {
return equityRequest -> {
ImdgValidationContext<FixedIncomeSecurityNewRequest> context = new ImdgValidationContext<>();
context.setValidatedObject(equityRequest);
context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, equityMap);
return new ValidatorImpl<>(context,
new NotPresentBySecuritySymbolAndWorkflowStatusActv(IMDGDistributedNames.Map_FixedIncomeSecurity, SecuritiesError.InstrumentAlreadyExists, FixedIncomeSecurityNewRequest.class),
new IsValidInstrumentTypeFromRequest(SecuritiesError.RequiredFieldIsEmpty, "BOND"),
FixedIncomeNewValidationRule.RequiredFieldIsNotEmpty
);
};
}
public Function<FixedIncomeSecurityUpdateRequest, IValidator> fixedIncomeUpdateValidator() {
return equityRequest -> {
ImdgValidationContext<FixedIncomeSecurityUpdateRequest> context = new ImdgValidationContext<>();
context.setValidatedObject(equityRequest);
context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, equityMap);
return new ValidatorImpl<>(context,
new PresentById(IMDGDistributedNames.Map_FixedIncomeSecurity, SecuritiesError.InstrumentNotFound, true),
new IsValidInstrumentTypeById(SecuritiesError.InstrumentNotFound, "BOND"),
new WorkflowStatusActv()
);
};
}
public Function<CommonDeleteRequest, IValidator> fixedIncomeDeleteValidator() {
return equityRequest -> {
ImdgValidationContext<CommonDeleteRequest> context = new ImdgValidationContext<>();
context.setValidatedObject(equityRequest);
context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, equityMap);
return new ValidatorImpl<>(context,
new PresentById(IMDGDistributedNames.Map_FixedIncomeSecurity, SecuritiesError.InstrumentNotFound, true),
new IsValidInstrumentTypeById(SecuritiesError.InstrumentNotFound, "BOND"),
new WorkflowStatusActv()
);
};

View file

@ -14,7 +14,7 @@ public enum EquityNewValidationRule implements IValidationRule<ImdgValidationCon
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<EquitySecurityNewRequest> context) {
EquitySecurityNewRequest action = context.getValidatedObject();
if (StringUtils.hasText(action.getShortName())){
if (!StringUtils.hasText(action.getShortName())){
return of(SecuritiesError.RequiredFieldIsEmpty, "shortName");
}
if (action.getLotSize() == null) {

View file

@ -0,0 +1,31 @@
package ru.spcex.clearing.securities.validation.rule;
import org.springframework.util.StringUtils;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.FixedIncomeSecurityNewRequest;
import ru.spcex.clearing.securities.errors.SecuritiesError;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Optional;
public enum FixedIncomeNewValidationRule implements IValidationRule<ImdgValidationContext<FixedIncomeSecurityNewRequest>> {
RequiredFieldIsNotEmpty() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<FixedIncomeSecurityNewRequest> context) {
FixedIncomeSecurityNewRequest action = context.getValidatedObject();
if (!StringUtils.hasText(action.getShortName())){
return of(SecuritiesError.RequiredFieldIsEmpty, "shortName");
}
if (action.getLotSize() == null) {
return of(SecuritiesError.RequiredFieldIsEmpty, "lotSize");
}
return empty();
}
};
@Override
public String ruleName() {
return "MmsNewValidationRule." + name();
}
}

View file

@ -0,0 +1,33 @@
package ru.spcex.clearing.securities.validation.rule;
import ru.spcex.clearing.securities.errors.SecuritiesError;
import ru.spcex.platform.classes.base.interfaces.WithInstrumentType;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.imdg.validation.Stored;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Optional;
public class IsValidInstrumentTypeById implements IValidationRule<ImdgValidationContext<WithInstrumentType>> {
private final IEnumId errorEnum;
private final String predictableInstrumentType;
public IsValidInstrumentTypeById(IEnumId errorEnum, String predictableInstrumentType) {
this.errorEnum = errorEnum;
this.predictableInstrumentType = predictableInstrumentType;
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<WithInstrumentType> context) {
WithInstrumentType action = context.getStoredObject(Stored.PresentById);
if (action.getInstrumentType() == null) {
return of(SecuritiesError.RequiredFieldIsEmpty);
} else if (!predictableInstrumentType.equalsIgnoreCase(action.getInstrumentType())){
return of(errorEnum);
}else {
return empty();
}
}
}

View file

@ -9,11 +9,11 @@ import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Optional;
public class IsValidInstrumentType implements IValidationRule<ImdgValidationContext<WithInstrumentType>> {
public class IsValidInstrumentTypeFromRequest implements IValidationRule<ImdgValidationContext<WithInstrumentType>> {
private final IEnumId errorEnum;
private final String predictableInstrumentType;
public IsValidInstrumentType(IEnumId errorEnum, String predictableInstrumentType) {
public IsValidInstrumentTypeFromRequest(IEnumId errorEnum, String predictableInstrumentType) {
this.errorEnum = errorEnum;
this.predictableInstrumentType = predictableInstrumentType;
}

View file

@ -14,7 +14,7 @@ public enum MmsNewValidationRule implements IValidationRule<ImdgValidationContex
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<MoneyMarketSecurityNewRequest> context) {
MoneyMarketSecurityNewRequest action = context.getValidatedObject();
if (StringUtils.hasText(action.getShortName())){
if (!StringUtils.hasText(action.getShortName())){
return of(SecuritiesError.RequiredFieldIsEmpty, "shortName");
}
if (action.getLotSize() == null) {
@ -23,10 +23,10 @@ public enum MmsNewValidationRule implements IValidationRule<ImdgValidationContex
if (action.getNominalValue() == null) {
return of(SecuritiesError.RequiredFieldIsEmpty, "nominalValue");
}
if (StringUtils.hasText(action.getNominalCurrency())) {
if (!StringUtils.hasText(action.getNominalCurrency())) {
return of(SecuritiesError.RequiredFieldIsEmpty, "nominalCurrency");
}
if (StringUtils.hasText(action.getTermType())) {
if (!StringUtils.hasText(action.getTermType())) {
return of(SecuritiesError.RequiredFieldIsEmpty, "termType");
}
return empty();

View file

@ -29,7 +29,7 @@ public class NotPresentBySecuritySymbolAndWorkflowStatusActv<T extends SpcexObje
public Optional<EnumMessage> validate(ImdgValidationContext<WithSecuritySymbol> context) {
WithSecuritySymbol action = context.getValidatedObject();
Imdg<T> imdg = context.obtainMap(mapName, clazz);
if(StringUtils.hasText(action.getSecuritySymbol())) {
if(!StringUtils.hasText(action.getSecuritySymbol())) {
return of(SecuritiesError.RequiredFieldIsEmpty, "securitySymbole");
}
T object = imdg.getSingleObjectByFieldValues(Map.of(

View file

@ -1,16 +0,0 @@
package ru.spcex.clearing.securities.config;
import com.hazelcast.client.HazelcastClient;
import com.hazelcast.core.HazelcastInstance;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class HazelcastInstanceTestConfiguration {
@Bean(name = "hazelcastInstance")
public HazelcastInstance hazelcastInstance() {
return HazelcastClient.newHazelcastClient();
}
}

View file

@ -3,10 +3,12 @@ package ru.spcex.clearing.securities.config;
import com.hazelcast.config.*;
import com.hazelcast.core.Hazelcast;
import com.hazelcast.core.HazelcastInstance;
import org.springframework.beans.factory.annotation.Autowired;
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.api.ImdgProvider;
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;
@ -15,10 +17,11 @@ import java.util.List;
import java.util.Random;
@Configuration
public class HazelcastServiceTestConfiguration {
public class ImdgTestConfig {
private HazelcastInstance hazelcastInstance;
private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
private static ThreadPoolTaskExecutor createThreadPoolTestTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
if (maxPoolSz > 2) {
pool.setKeepAliveSeconds(60);
@ -29,34 +32,34 @@ public class HazelcastServiceTestConfiguration {
return pool;
}
@Bean(name = "taskExecutorHazelcastTestClientInitializer")
public ThreadPoolTaskExecutor taskExecutorHazelcastTestClientInitializer() {
return createThreadPoolTestTaskExecutor(1, true);
}
@Bean(name = "taskExecutorTestIdGeneratorAwaiter")
public ThreadPoolTaskExecutor taskExecutorTestIdGeneratorAwaiter() {
return createThreadPoolTestTaskExecutor(1, false);
}
@Autowired
@Bean(name = "hazelcastServiceTest")
public HazelcastService hazelcastService(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, HazelcastClientParams params) {
public ImdgProvider imdgTestProvider(
@Qualifier("taskExecutorHazelcastTestClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
@Qualifier("taskExecutorTestIdGeneratorAwaiter") 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")));
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);
}
@Bean(name = "taskExecutorHazelcastClientInitializer")
public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
return createThreadPoolTaskExecutor(1, true);
}
@Bean(name = "taskExecutorIdGeneratorAwaiter")
public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
return createThreadPoolTaskExecutor(1, false);
return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params);
}
@Bean(name = "hazelcastClientParams")
@ -67,7 +70,6 @@ public class HazelcastServiceTestConfiguration {
params.setClusterMembers("127.0.0.1");
params.setInstanceName("hzTestClient" + new Random().nextInt());
params.setNearCacheConfig(new NearCacheConfig());
return params;
}
}
}

View file

@ -0,0 +1,42 @@
package ru.spcex.clearing.securities.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.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 KafkaTestConfig {
@Autowired
@Bean
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer,
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
public MockConsumer<String, Object> createTestConsumer() {
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
}
}

View file

@ -0,0 +1,74 @@
package ru.spcex.clearing.securities.service;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Captor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.securities.component.ListingBuilder;
import ru.spcex.clearing.securities.config.ErrorResolverConfig;
import ru.spcex.clearing.securities.config.ImdgTestConfig;
import ru.spcex.clearing.securities.config.KafkaTestConfig;
import ru.spcex.clearing.securities.config.ValidationConfig;
import ru.spcex.clearing.securities.service.cud.EquitySecurityService;
import ru.spcex.clearing.securities.service.cud.FixedIncomeSecurityService;
import ru.spcex.clearing.securities.service.cud.MoneyMarketSecurityService;
import ru.spcex.clearing.securities.utils.MatcherFactory;
import ru.spcex.clearing.securities.utils.TestUtils;
import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.concurrent.atomic.AtomicLong;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static ru.spcex.clearing.securities.utils.MatcherFactory.usingIgnoringFieldsComparator;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
EquitySecurityService.class,
FixedIncomeSecurityService.class,
MoneyMarketSecurityService.class,
ImdgTestConfig.class,
KafkaTestConfig.class,
ValidationProvider.class,
ErrorResolverConfig.class,
ValidationConfig.class})
public abstract class AbstractServiceTest {
protected static final MatcherFactory.Matcher<Listing> LISTING_MATCHER = usingIgnoringFieldsComparator("created", "updated");
protected static AtomicLong currentId = new AtomicLong(0L);
protected Imdg<MoneyMarketSecurity> moneyMarketSecurityMap;
protected Imdg<Listing> listingImdg;
protected Imdg<EquitySecurity> equitySecurityImdg;
protected Imdg<FixedIncomeSecurity> fixedIncomeSecurityImdg;
protected final ListingBuilder listingBuilder = new ListingBuilder();
@Captor
protected ArgumentCaptor<ProducerRecord> producerRecord;
@MockBean
protected MockProducer<String, Object> mockProducer;
@Autowired
@Qualifier("hazelcastServiceTest")
protected ImdgProvider imdgProvider;
protected void init() {
imdgProvider.waitAvailable();
this.moneyMarketSecurityMap = imdgProvider.getImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, MoneyMarketSecurity.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class);
this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
TestUtils.FutureRecordMetadata future = spy(TestUtils.FutureRecordMetadata.class);
doReturn(future).when(mockProducer).send(producerRecord.capture());
}
}

View file

@ -0,0 +1,210 @@
package ru.spcex.clearing.securities.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
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.securitites.EquitySecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest;
import ru.spcex.clearing.securities.service.cud.EquitySecurityService;
import ru.spcex.clearing.securities.service.cud.MoneyMarketSecurityService;
import ru.spcex.clearing.securities.utils.MatcherFactory;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.util.Map;
import static ru.spcex.clearing.securities.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.securities.utils.TestUtils.*;
public class EquitySecurityServiceTest extends AbstractServiceTest{
private static final MatcherFactory.Matcher<EquitySecurity> EQUITY_SECURITY_MATCHER = usingIgnoringFieldsComparator("created", "updated");
private final long ID = currentId.getAndIncrement();
private final int PARTITION = 0;
@Autowired
private EquitySecurityService equitySecurityService;
@PostConstruct
protected void init() {
super.init();
}
/**
* {@link MoneyMarketSecurityService#newMoneyMarket(BaseRequest)}<br>
* Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityNewRequest}:<br>
* {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L<br>
* {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"<br>
* {@link MoneyMarketSecurityNewRequest#fullName} - "estat"<br>
*/
@Test
public void testNewEquity(){
final String TOPIC = Consts.DESTINATION_EQUITY_SECURITY_NEW;
final EquitySecurity equityPrediction = getEquitySecurity();
final EquitySecurityNewRequest request = getEquitySecurityNewRequest(equityPrediction);
Listing listingPrediction = listingBuilder.byNewEquity(equityPrediction, request);
//ACT
String jsonString = getJsonStringForNew(request, ID);
addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
equityPrediction.setId(equityResult.getId());
equityPrediction.setSecurityId(equityResult.getSecurityId());
EQUITY_SECURITY_MATCHER.assertMatch(equityResult, equityPrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equityPrediction.getId()));
listingPrediction.setId(listingResult.getId());
listingPrediction.setSecurityId(listingResult.getSecurityId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
/**
* {@link MoneyMarketSecurityService#updateMoneyMarket(BaseRequest)}
* Тест проверяет обновление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityUpdateRequest}:<br>
* {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg<br>
* {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0<br>
* {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)<br>
*/
@Test
public void testUpdateEquity() {
final String TOPIC = Consts.DESTINATION_EQUITY_SECURITY_UPDATE;
String fullName = "FullNameUpdate";
Double updateLotSize = 12.0;
final EquitySecurity equityPrediction = getEquitySecurity();
equitySecurityImdg.insert(equityPrediction);
equityPrediction.setFullName(fullName);
equityPrediction.setLotSize(BigDecimal.valueOf(updateLotSize));
final EquitySecurityUpdateRequest request = getEquitySecurityUpdateRequest(equityPrediction);
Listing listingPrediction = new Listing();
listingPrediction = listingBuilder.byUpdateEquity(equityPrediction, listingPrediction);
listingImdg.insert(listingPrediction);
//ACT
String jsonString = getJsonStringForUPDATE(request, ID);
addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
equityPrediction.setId(equityResult.getId());
equityPrediction.setSecurityId(equityResult.getSecurityId());
EQUITY_SECURITY_MATCHER.assertMatch(equityResult, equityPrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equityPrediction.getId()));
listingPrediction.setId(listingResult.getId());
listingPrediction.setSecurityId(listingResult.getSecurityId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
/**
* {@link MoneyMarketSecurityService#deleteMoneyMarket(BaseRequest)}<br>
* Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link CommonDeleteRequest}:<br>
* {@link CommonDeleteRequest#id} - Идентификатор записи<br>
*/
@Test
public void testDeleteEquity(){
final String TOPIC = Consts.DESTINATION_EQUITY_SECURITY_DELETE;
final EquitySecurity equityPrediction = getEquitySecurity();
equitySecurityImdg.insert(equityPrediction);
equityPrediction.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
Listing listingPrediction = new Listing();
listingPrediction = listingBuilder.byUpdateEquity(equityPrediction, listingPrediction);
listingImdg.insert(listingPrediction);listingPrediction.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
commonDeleteRequest.setId(ID);
//ACT
String jsonString = getJsonStringForDELETE(commonDeleteRequest, ID);
addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
equityPrediction.setId(equityResult.getId());
equityPrediction.setSecurityId(equityResult.getSecurityId());
EQUITY_SECURITY_MATCHER.assertMatch(equityResult, equityPrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equityPrediction.getId()));
listingPrediction.setId(listingResult.getId());
listingPrediction.setSecurityId(listingResult.getSecurityId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
private EquitySecurity getEquitySecurity(){
EquitySecurity equitySecurity = new EquitySecurity();
equitySecurity.setId(ID);
equitySecurity.setSecuritySymbol("SecuritySymbol");
equitySecurity.setShortName("ShortNameNewEquity");
equitySecurity.setFullName("FullNameNewEquity");
equitySecurity.setIsin("Isin");
equitySecurity.setShareType("ShareType");
equitySecurity.setLotSize(BigDecimal.valueOf(18.0d));
equitySecurity.setIssuerId(11L);
equitySecurity.setShortNameEng("ShortNameNewEquity");
equitySecurity.setFullNameEng("FullNameNewEquity");
equitySecurity.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
equitySecurity.setInstrumentType("EQTY");
return equitySecurity;
}
private EquitySecurityNewRequest getEquitySecurityNewRequest(EquitySecurity equitySecurity){
EquitySecurityNewRequest equitySecurityNewRequest = new EquitySecurityNewRequest();
equitySecurityNewRequest.setSecuritySymbol(equitySecurity.getSecuritySymbol());
equitySecurityNewRequest.setShortName(equitySecurity.getShortName());
equitySecurityNewRequest.setFullName(equitySecurity.getFullName());
equitySecurityNewRequest.setIsin(equitySecurity.getIsin());
equitySecurityNewRequest.setShareType(equitySecurity.getShareType());
equitySecurityNewRequest.setLotSize(equitySecurity.getLotSize().doubleValue());
equitySecurityNewRequest.setIssuerId(equitySecurity.getIssuerId());
equitySecurityNewRequest.setShortNameEng(equitySecurity.getShortNameEng());
equitySecurityNewRequest.setFullNameEng(equitySecurity.getFullNameEng());
equitySecurityNewRequest.setWorkflowStatus(equitySecurity.getWorkflowStatus());
equitySecurityNewRequest.setInstrumentType(equitySecurity.getInstrumentType());
return equitySecurityNewRequest;
}
private EquitySecurityUpdateRequest getEquitySecurityUpdateRequest(EquitySecurity equitySecurity){
EquitySecurityUpdateRequest equitySecurityUpdateRequest = new EquitySecurityUpdateRequest();
equitySecurityUpdateRequest.setId(equitySecurity.getId());
equitySecurityUpdateRequest.setSecuritySymbol(equitySecurity.getSecuritySymbol());
equitySecurityUpdateRequest.setShortName(equitySecurity.getShortName());
equitySecurityUpdateRequest.setFullName(equitySecurity.getFullName());
equitySecurityUpdateRequest.setIsin(equitySecurity.getIsin());
equitySecurityUpdateRequest.setShareType(equitySecurity.getShareType());
equitySecurityUpdateRequest.setLotSize(equitySecurity.getLotSize().doubleValue());
equitySecurityUpdateRequest.setIssuerId(equitySecurity.getIssuerId());
equitySecurityUpdateRequest.setShortNameEng(equitySecurity.getShortNameEng());
equitySecurityUpdateRequest.setFullNameEng(equitySecurity.getFullNameEng());
equitySecurityUpdateRequest.setWorkflowStatus(equitySecurity.getWorkflowStatus());
equitySecurityUpdateRequest.setInstrumentType(equitySecurity.getInstrumentType());
return equitySecurityUpdateRequest;
}
}

View file

@ -0,0 +1,227 @@
package ru.spcex.clearing.securities.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
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.securitites.EquitySecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.FixedIncomeSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest;
import ru.spcex.clearing.securities.service.cud.FixedIncomeSecurityService;
import ru.spcex.clearing.securities.service.cud.MoneyMarketSecurityService;
import ru.spcex.clearing.securities.utils.MatcherFactory;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.util.Map;
import static ru.spcex.clearing.securities.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.securities.utils.TestUtils.*;
public class FixedIncomeSecurityServiceTest extends AbstractServiceTest{
private static final MatcherFactory.Matcher<FixedIncomeSecurity> FIXED_INCOME_SECURITY_MATCHER = usingIgnoringFieldsComparator("created", "updated");
private final long ID = currentId.getAndIncrement();
private final int PARTITION = 0;
@Autowired
private FixedIncomeSecurityService fixedIncomeSecurityService;
@PostConstruct
protected void init() {
super.init();
}
/**
* {@link MoneyMarketSecurityService#newMoneyMarket(BaseRequest)}<br>
* Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityNewRequest}:<br>
* {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L<br>
* {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"<br>
* {@link MoneyMarketSecurityNewRequest#fullName} - "estat"<br>
*/
@Test
public void newFixedIncome(){
final String TOPIC = Consts.DESTINATION_FIXED_INCOME_SECURITY_NEW;
final FixedIncomeSecurity fixedIncomePrediction = getFixedIncomeSecurity();
final FixedIncomeSecurityNewRequest request = getFixedIncomeSecurityNewRequest(fixedIncomePrediction);
Listing listingPrediction = listingBuilder.byNewFixedIncome(fixedIncomePrediction, request);
//ACT
String jsonString = getJsonStringForNew(request, ID);
addRecordToKafka((MockConsumer) fixedIncomeSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
FixedIncomeSecurity equityResult = fixedIncomeSecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName()));
fixedIncomePrediction.setId(equityResult.getId());
fixedIncomePrediction.setSecurityId(equityResult.getSecurityId());
FIXED_INCOME_SECURITY_MATCHER.assertMatch(equityResult, fixedIncomePrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", fixedIncomePrediction.getId()));
listingPrediction.setId(listingResult.getId());
listingPrediction.setSecurityId(listingResult.getSecurityId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
/**
* {@link MoneyMarketSecurityService#updateMoneyMarket(BaseRequest)}
* Тест проверяет обновление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityUpdateRequest}:<br>
* {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg<br>
* {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0<br>
* {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)<br>
*/
@Test
public void updateFixedIncome() {
// clearAllInImdg(moneyMarketSecurityMap);
// final String TOPIC = Consts.DESTINATION_FIXED_INCOME_SECURITY_UPDATE;
// String fullName = "FullNameUpdate";
// Double updateLotSize = 12.0;
//
// final FixedIncomeSecurity fixedIncomePrediction = getFixedIncomeSecurity();
//
// fixedIncomeSecurityImdg.insert(fixedIncomePrediction);
// fixedIncomePrediction.setFullName(fullName);
// fixedIncomePrediction.setLotSize(BigDecimal.valueOf(updateLotSize));
// final FixedIncomeSecurityUpdateRequest request = getEquitySecurityUpdateRequest(fixedIncomePrediction);
// Listing listingPrediction = new Listing();
// listingPrediction = listingBuilder.byUpdateEquity(fixedIncomePrediction, listingPrediction);
// listingImdg.insert(listingPrediction);
//
// //ACT
// String jsonString = getJsonStringForUPDATE(request, ID);
//
// addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//
// //ASSERT
// waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
//
// EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", fixedIncomePrediction.getFullName()));
// fixedIncomePrediction.setId(equityResult.getId());
// fixedIncomePrediction.setSecurityId(equityResult.getSecurityId());
// FIXED_INCOME_SECURITY_MATCHER.assertMatch(equityResult, fixedIncomePrediction);
//
// Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", fixedIncomePrediction.getId()));
// listingPrediction.setId(listingResult.getId());
// listingPrediction.setSecurityId(listingResult.getSecurityId());
// LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
/**
* {@link MoneyMarketSecurityService#deleteMoneyMarket(BaseRequest)}<br>
* Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link CommonDeleteRequest}:<br>
* {@link CommonDeleteRequest#id} - Идентификатор записи<br>
*/
@Test
public void deleteFixedIncome(){
// clearAllInImdg(moneyMarketSecurityMap);
// final String TOPIC = Consts.DESTINATION_FIXED_INCOME_SECURITY_DELETE;
//
// final EquitySecurity equityPrediction = getFixedIncomeSecurity();
// equitySecurityImdg.insert(equityPrediction);
// equityPrediction.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
//
// Listing listingPrediction = new Listing();
// listingPrediction = listingBuilder.byUpdateEquity(equityPrediction, listingPrediction);
// listingImdg.insert(listingPrediction);listingPrediction.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
//
// CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
// commonDeleteRequest.setId(ID);
//
// //ACT
// String jsonString = getJsonStringForDELETE(commonDeleteRequest, ID);
//
// addRecordToKafka((MockConsumer) equitySecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//
// //ASSERT
// waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
//
// EquitySecurity equityResult = equitySecurityImdg.getSingleObjectBySQL(String.format("fullName = %s", equityPrediction.getFullName()));
// equityPrediction.setId(equityResult.getId());
// equityPrediction.setSecurityId(equityResult.getSecurityId());
// FIXED_INCOME_SECURITY_MATCHER.assertMatch(equityResult, equityPrediction);
//
// Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", equityPrediction.getId()));
// listingPrediction.setId(listingResult.getId());
// listingPrediction.setSecurityId(listingResult.getSecurityId());
// LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
private FixedIncomeSecurity getFixedIncomeSecurity(){
FixedIncomeSecurity fixedIncome = new FixedIncomeSecurity();
fixedIncome.setId(ID);
fixedIncome.setSecuritySymbol("SecuritySymbol");
fixedIncome.setShortName("ShortNameNewEquity");
fixedIncome.setFullName("FullNameNewEquity");
fixedIncome.setIsin("Isin");
fixedIncome.setBondType("BondType");
fixedIncome.setLotSize(BigDecimal.valueOf(18.0d));
fixedIncome.setNominalValue(BigDecimal.valueOf(18.0d));
fixedIncome.setNominalCurrency("NominalCurrency");
fixedIncome.setMaturityDate(LocalDate.now());
fixedIncome.setCoupon(BigDecimal.valueOf(18.0d));
fixedIncome.setCouponFrequency(9L);
fixedIncome.setIssuerId(11L);
fixedIncome.setShortNameEng("ShortNameNewEquity");
fixedIncome.setFullNameEng("FullNameNewEquity");
fixedIncome.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
fixedIncome.setInstrumentType("BOND");
fixedIncome.setSecurityId(fixedIncome.getId());
return fixedIncome;
}
private FixedIncomeSecurityNewRequest getFixedIncomeSecurityNewRequest(FixedIncomeSecurity fixedIncomeSecurity){
FixedIncomeSecurityNewRequest equitySecurityNewRequest = new FixedIncomeSecurityNewRequest();
equitySecurityNewRequest.setSecuritySymbol(fixedIncomeSecurity.getSecuritySymbol());
equitySecurityNewRequest.setShortName(fixedIncomeSecurity.getShortName());
equitySecurityNewRequest.setFullName(fixedIncomeSecurity.getFullName());
equitySecurityNewRequest.setIsin(fixedIncomeSecurity.getIsin());
equitySecurityNewRequest.setBondType(fixedIncomeSecurity.getBondType());
equitySecurityNewRequest.setLotSize(fixedIncomeSecurity.getLotSize().doubleValue());
equitySecurityNewRequest.setNominalValue(fixedIncomeSecurity.getNominalValue().doubleValue());
equitySecurityNewRequest.setNominalCurrency(fixedIncomeSecurity.getNominalCurrency());
equitySecurityNewRequest.setMaturityDate(fixedIncomeSecurity.getMaturityDate());
equitySecurityNewRequest.setCoupon(fixedIncomeSecurity.getCoupon().doubleValue());
equitySecurityNewRequest.setCouponFrequency(9L);
equitySecurityNewRequest.setIssuerId(fixedIncomeSecurity.getIssuerId());
equitySecurityNewRequest.setShortNameEng(fixedIncomeSecurity.getShortNameEng());
equitySecurityNewRequest.setFullNameEng(fixedIncomeSecurity.getFullNameEng());
equitySecurityNewRequest.setWorkflowStatus(fixedIncomeSecurity.getWorkflowStatus());
equitySecurityNewRequest.setInstrumentType(fixedIncomeSecurity.getInstrumentType());
return equitySecurityNewRequest;
}
private EquitySecurityUpdateRequest getEquitySecurityUpdateRequest(EquitySecurity equitySecurity){
EquitySecurityUpdateRequest equitySecurityUpdateRequest = new EquitySecurityUpdateRequest();
equitySecurityUpdateRequest.setId(equitySecurity.getId());
equitySecurityUpdateRequest.setSecuritySymbol(equitySecurity.getSecuritySymbol());
equitySecurityUpdateRequest.setShortName(equitySecurity.getShortName());
equitySecurityUpdateRequest.setFullName(equitySecurity.getFullName());
equitySecurityUpdateRequest.setIsin(equitySecurity.getIsin());
equitySecurityUpdateRequest.setShareType(equitySecurity.getShareType());
equitySecurityUpdateRequest.setLotSize(equitySecurity.getLotSize().doubleValue());
equitySecurityUpdateRequest.setIssuerId(equitySecurity.getIssuerId());
equitySecurityUpdateRequest.setShortNameEng(equitySecurity.getShortNameEng());
equitySecurityUpdateRequest.setFullNameEng(equitySecurity.getFullNameEng());
equitySecurityUpdateRequest.setWorkflowStatus(equitySecurity.getWorkflowStatus());
equitySecurityUpdateRequest.setInstrumentType(equitySecurity.getInstrumentType());
return equitySecurityUpdateRequest;
}
}

View file

@ -16,11 +16,17 @@ class MoneyMarketSecurityFactory {
private String description = "tsetse";
private Double lotSize = 1.0;
private LocalDate startDate = LocalDate.ofEpochDay(0);
private LocalDate endDate = LocalDate.ofEpochDay(0);
private LocalDate endDate = LocalDate.now();
private BigDecimal nominalValue = BigDecimal.valueOf(1.0);
private String instrumentType = "sdgaga"; // (linked to instrumentType)
private String instrumentType = "RATE"; // (linked to instrumentType)
private String fullName = "estat";
private String shortName = "esasdftat";
private String securitySymbol = "sdfafd";
private String termType = "termType";
public String getFullName() {
return fullName;
}
public MoneyMarketSecurityFactory setLotSize(Double lotSize) {
this.lotSize = lotSize;
@ -89,7 +95,12 @@ class MoneyMarketSecurityFactory {
moneyMarketSecurity.setNominalCurrency(nominalCurrency);
moneyMarketSecurity.setInstrumentType(instrumentType);
moneyMarketSecurity.setFullName(fullName);
moneyMarketSecurity.setShortName(shortName);
moneyMarketSecurity.setTermType(termType);
moneyMarketSecurity.setDescription(description);
moneyMarketSecurity.setLotSize(BigDecimal.valueOf(lotSize));
moneyMarketSecurity.setSecuritySymbol(securitySymbol);
moneyMarketSecurity.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Active.getKey());
return moneyMarketSecurity;
}
@ -99,11 +110,17 @@ class MoneyMarketSecurityFactory {
= new MoneyMarketSecurityUpdateRequest();
moneyMarketSecurityUpdateRequest.setId(id);
moneyMarketSecurityUpdateRequest.setLotSize(lotSize);
moneyMarketSecurityUpdateRequest.setStartDate(startDate);
moneyMarketSecurityUpdateRequest.setEndDate(endDate);
moneyMarketSecurityUpdateRequest.setNominalValue(nominalValue.doubleValue());
moneyMarketSecurityUpdateRequest.setNominalCurrency(nominalCurrency);
moneyMarketSecurityUpdateRequest.setInstrumentType(instrumentType);
moneyMarketSecurityUpdateRequest.setFullName(fullName);
moneyMarketSecurityUpdateRequest.setShortName(shortName);
moneyMarketSecurityUpdateRequest.setTermType(termType);
moneyMarketSecurityUpdateRequest.setDescription(description);
moneyMarketSecurityUpdateRequest.setSecuritySymbol(securitySymbol);
return moneyMarketSecurityUpdateRequest;
}
@ -124,6 +141,9 @@ class MoneyMarketSecurityFactory {
keyRequest.setNominalCurrency(nominalCurrency);
keyRequest.setInstrumentType(instrumentType);
keyRequest.setFullName(fullName);
keyRequest.setShortName(shortName);
keyRequest.setTermType(termType);
keyRequest.setDescription(description);
keyRequest.setSecuritySymbol(securitySymbol);
return keyRequest;

View file

@ -1,324 +0,0 @@
package ru.spcex.clearing.securities.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.IMap;
import com.hazelcast.map.listener.EntryAddedListener;
import com.hazelcast.map.listener.EntryRemovedListener;
import com.hazelcast.map.listener.EntryUpdatedListener;
import org.apache.kafka.clients.consumer.ConsumerRecord;
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.common.TopicPartition;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.misc.KeyRate;
import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
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.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest;
import ru.spcex.clearing.securities.config.HazelcastInstanceTestConfiguration;
import ru.spcex.clearing.securities.config.HazelcastServiceTestConfiguration;
import ru.spcex.clearing.securities.service.cud.MoneyMarketSecurityService;
import ru.spcex.clearing.securities.utils.MatcherFactory;
import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
import ru.spcex.platform.utils.enumeration.SimpleMessageResolver;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.util.Collections;
import java.util.HashMap;
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
HazelcastServiceTestConfiguration.class,
HazelcastInstanceTestConfiguration.class})
public class MoneyMarketSecurityServiceDisabled {
static final ObjectMapper objectMapper = new ObjectMapper();
private final int TIMEOUT = 1000;
boolean isUsed;
private MockConsumer<String, Object> kafkaMockQueue;
private final MockProducer<String, Object> mockProducer = new MockProducer<>();
// @Autowired
// @Qualifier("hazelcastServiceTest")
private HazelcastService hazelcastService;
// @Autowired
// @Qualifier("hazelcastInstance")
private HazelcastInstance hz;
private void setIsUsed(boolean isUsed) {
this.isUsed = isUsed;
}
private String jsonBaseRequest(BaseRequest baseRequest) {
String jsonBaseRequest;
try {
jsonBaseRequest = objectMapper.writeValueAsString(baseRequest);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
return jsonBaseRequest;
}
void initKafka(String topic, int partition, Long offset) {
TopicPartition tp = new TopicPartition(topic, partition);
kafkaMockQueue.schedulePollTask(() -> {
kafkaMockQueue.rebalance(Collections.singletonList(tp));
});
HashMap<TopicPartition, Long> startOffsets = new HashMap<>();
startOffsets.put(tp, offset);
kafkaMockQueue.updateBeginningOffsets(startOffsets);
}
private void startService() throws InterruptedException {
MoneyMarketSecurityService keyRateService =
new MoneyMarketSecurityService(kafkaMockQueue, mockProducer, hazelcastService,
new ValidationProvider(hazelcastService), new SimpleMessageResolver());
keyRateService.afterPropertiesSet();
}
/**
* {@link MoneyMarketSecurityService#newMoneyMarket(BaseRequest)}<br>
* Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityNewRequest}:<br>
* {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L<br>
* {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"<br>
* {@link MoneyMarketSecurityNewRequest#fullName} - "estat"<br>
*/
// @Test
public void testNewMoneyMarketSecurity() throws InterruptedException {
final String TOPIC = Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW;
final String MAP = IMDGDistributedNames.Map_MoneyMarketSecurity;
final int PARTITION = 0;
final Long OFFSET = 0L;
final Long REQUEST_ID = 0L;
final MoneyMarketSecurityFactory moneyMarketSecurityFactory = new MoneyMarketSecurityFactory();
final MoneyMarketSecurity moneyMarketSecurityPrediction = moneyMarketSecurityFactory
.getMoneyMarketSecurity();
final MoneyMarketSecurityNewRequest keyRequest = moneyMarketSecurityFactory
.getMoneyMarketSecurityNewRequest();
BaseRequest baseRequest = new BaseRequest();
baseRequest.setRequestPayload(keyRequest);
baseRequest.setId(REQUEST_ID);
baseRequest.setActionType(ActionType.NEW);
String jsonRequest = jsonBaseRequest(baseRequest);
ConsumerRecord<String, Object> consumerRecord =
new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonRequest);
kafkaMockQueue = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
initKafka(TOPIC, PARTITION, OFFSET);
kafkaMockQueue.schedulePollTask(() -> kafkaMockQueue.addRecord(consumerRecord));
isUsed = false;
IMap<Long, MoneyMarketSecurity> iMap = hz.getMap(MAP);
Object waiter = new Object();
String listenerID = iMap.addEntryListener((EntryAddedListener<Long, MoneyMarketSecurity>) entryEvent -> {
System.out.println("Checking Equality");
MoneyMarketSecurity moneyMarketSecurityResult = entryEvent.getValue();
moneyMarketSecurityPrediction.setId(entryEvent.getKey());
MatcherFactory.Matcher<MoneyMarketSecurity> matcher = MatcherFactory
.usingIgnoringFieldsComparator("id", "securityId", "description");
matcher.assertMatch(moneyMarketSecurityResult, moneyMarketSecurityPrediction);
System.out.println("ok");
setIsUsed(true);
synchronized (waiter) {
waiter.notify();
}
}, true);
startService();
synchronized (waiter) {
waiter.wait(TIMEOUT);
}
Assertions.assertEquals(true, isUsed);
iMap.removeEntryListener(listenerID);
}
/**
* {@link MoneyMarketSecurityService#deleteMoneyMarket(BaseRequest)}<br>
* Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link CommonDeleteRequest}:<br>
* {@link CommonDeleteRequest#id} - Идентификатор записи<br>
*/
// @Test
public void testDeleteKeyRate() throws InterruptedException {
final String TOPIC = Consts.DESTINATION_MONEY_MARKET_SECURITY_DELETE;
final String MAP = IMDGDistributedNames.Map_MoneyMarketSecurity;
final int PARTITION = 0;
final Long ID = 1L;
final Long OFFSET = 0L;
final Long REQUEST_ID = 2L;
final MoneyMarketSecurityFactory moneyMarketSecurityFactory =
new MoneyMarketSecurityFactory()
.setId(ID);
final MoneyMarketSecurity initialMoneyMarketSecurity = moneyMarketSecurityFactory
.getMoneyMarketSecurity();
final CommonDeleteRequest keyRequest = moneyMarketSecurityFactory
.getCommonDeleteRequest();
BaseRequest baseRequest = new BaseRequest();
baseRequest.setRequestPayload(keyRequest);
baseRequest.setId(REQUEST_ID);
baseRequest.setActionType(ActionType.DELETE);
String jsonRequest = jsonBaseRequest(baseRequest);
ConsumerRecord<String, Object> consumerRecord =
new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonRequest);
kafkaMockQueue = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
initKafka(TOPIC, PARTITION, OFFSET);
kafkaMockQueue.schedulePollTask(() -> kafkaMockQueue.addRecord(consumerRecord));
isUsed = false;
IMap<Long, MoneyMarketSecurity> iMap = hz.getMap(MAP);
iMap.put(ID, initialMoneyMarketSecurity);
Object waiter = new Object();
String listenerID = iMap.addEntryListener((EntryRemovedListener<Long, KeyRate>) entryEvent -> {
System.out.println("Checking If removed..");
Assertions.assertEquals(ID, entryEvent.getKey());
System.out.println("ok");
setIsUsed(true);
synchronized (waiter) {
waiter.notify();
}
}, false);
startService();
synchronized (waiter) {
waiter.wait(TIMEOUT);
}
Assertions.assertEquals(true, isUsed);
iMap.removeEntryListener(listenerID);
}
/**
* {@link MoneyMarketSecurityService#updateMoneyMarket(BaseRequest)}
* Тест проверяет обновление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityUpdateRequest}:<br>
* {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg<br>
* {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0<br>
* {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)<br>
*/
// @Test
public void testUpdateKeyRate() throws InterruptedException {
final String TOPIC = Consts.DESTINATION_MONEY_MARKET_SECURITY_UPDATE;
final String MAP = IMDGDistributedNames.Map_MoneyMarketSecurity;
final int PARTITION = 0;
final Long ID = 1L;
final Long OFFSET = 0L;
final Long REQUEST_ID = 1L;
final MoneyMarketSecurityFactory moneyMarketSecurityFactory =
new MoneyMarketSecurityFactory()
.setId(ID);
final MoneyMarketSecurity initialMoneyMarketSecurity = moneyMarketSecurityFactory
.getMoneyMarketSecurity();
moneyMarketSecurityFactory
.setFullName("sfgsdfg")
.setInstrumentType("qwerqwe")
.setNominalCurrency("RUB")
.setNominalValue(BigDecimal.valueOf(2.0))
.setEndDate(LocalDate.ofEpochDay(24 * 3600));
final MoneyMarketSecurity moneyMarketSecurityPrediction =
moneyMarketSecurityFactory.getMoneyMarketSecurity();
final MoneyMarketSecurityUpdateRequest keyRequest =
moneyMarketSecurityFactory.getMoneyMarketSecurityUpdateRequest();
BaseRequest baseRequest = new BaseRequest();
baseRequest.setRequestPayload(keyRequest);
baseRequest.setId(REQUEST_ID);
baseRequest.setActionType(ActionType.UPDATE);
String jsonRequest = jsonBaseRequest(baseRequest);
ConsumerRecord<String, Object> consumerRecord =
new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonRequest);
kafkaMockQueue = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
initKafka(TOPIC, PARTITION, OFFSET);
kafkaMockQueue.schedulePollTask(() -> kafkaMockQueue.addRecord(consumerRecord));
IMap<Long, MoneyMarketSecurity> iMap = hz.getMap(MAP);
iMap.put(ID, initialMoneyMarketSecurity);
isUsed = false;
Object waiter = new Object();
String listenerID = iMap.addEntryListener((EntryUpdatedListener<Long, MoneyMarketSecurity>) entryEvent -> {
System.out.println("Checking Equality..");
MoneyMarketSecurity moneyMarketSecurityResult = entryEvent.getValue();
MatcherFactory.Matcher<MoneyMarketSecurity> matcher = MatcherFactory
.usingIgnoringFieldsComparator("id");
matcher.assertMatch(moneyMarketSecurityResult, moneyMarketSecurityPrediction);
System.out.println("Ok");
setIsUsed(true);
synchronized (waiter) {
waiter.notify();
}
}, true);
startService();
synchronized (waiter) {
waiter.wait(TIMEOUT);
}
Assertions.assertEquals(true, isUsed);
iMap.removeEntryListener(listenerID);
}
}

View file

@ -0,0 +1,179 @@
package ru.spcex.clearing.securities.service;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity;
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.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest;
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
import ru.spcex.clearing.platform.messaging.service.Status;
import ru.spcex.clearing.securities.service.cud.MoneyMarketSecurityService;
import ru.spcex.clearing.securities.utils.MatcherFactory;
import ru.spcex.platform.imdg.api.Imdg;
import javax.annotation.PostConstruct;
import java.math.BigDecimal;
import java.util.Map;
import static ru.spcex.clearing.securities.utils.MatcherFactory.usingIgnoringFieldsComparator;
import static ru.spcex.clearing.securities.utils.TestUtils.*;
public class MoneyMarketSecurityServiceTest extends AbstractServiceTest{
private static final MatcherFactory.Matcher<MoneyMarketSecurity> MONEY_MARKET_SECURITY_MATCHER = usingIgnoringFieldsComparator("created", "updated");
private final long ID = currentId.getAndIncrement();
private final int PARTITION = 0;
private MoneyMarketSecurityFactory moneyMarketSecurityFactory;
@Autowired
private MoneyMarketSecurityService moneyMarketSecurityService;
@PostConstruct
protected void init() {
super.init();
Imdg<RequestInfo> requestSpecificImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
RequestInfo requestFound = new RequestInfo();
requestFound.setId(ID);
requestFound.setStatus(Status.Processing);
requestSpecificImdg.insert(requestFound);
moneyMarketSecurityFactory = new MoneyMarketSecurityFactory();
moneyMarketSecurityFactory.setId(ID);
}
/**
* {@link MoneyMarketSecurityService#newMoneyMarket(BaseRequest)}<br>
* Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityNewRequest}:<br>
* {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.<br>
* {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0<br>
* {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L<br>
* {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"<br>
* {@link MoneyMarketSecurityNewRequest#fullName} - "estat"<br>
*/
@Test
public void testNewMoneyMarketSecurity(){
clearAllInImdg(moneyMarketSecurityMap);
final String TOPIC = Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW;
final MoneyMarketSecurity moneyMarketSecurityPrediction = moneyMarketSecurityFactory
.getMoneyMarketSecurity();
final MoneyMarketSecurityNewRequest keyRequest = moneyMarketSecurityFactory
.getMoneyMarketSecurityNewRequest();
Listing listingPrediction = listingBuilder.byNewMms(moneyMarketSecurityPrediction, keyRequest);
//ACT
String jsonString = getJsonStringForNew(keyRequest, ID);
addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName()));
moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId());
moneyMarketSecurityPrediction.setSecurityId(moneyMarketSecurityResult.getSecurityId());
MONEY_MARKET_SECURITY_MATCHER.assertMatch(moneyMarketSecurityResult, moneyMarketSecurityPrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", moneyMarketSecurityPrediction.getId()));
listingPrediction.setId(listingResult.getId());
listingPrediction.setSecurityId(listingResult.getSecurityId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
/**
* {@link MoneyMarketSecurityService#deleteMoneyMarket(BaseRequest)}<br>
* Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link CommonDeleteRequest}:<br>
* {@link CommonDeleteRequest#id} - Идентификатор записи<br>
*/
@Test
public void testDeleteMoneyMarket(){
clearAllInImdg(moneyMarketSecurityMap);
final String TOPIC = Consts.DESTINATION_MONEY_MARKET_SECURITY_DELETE;
final MoneyMarketSecurity moneyMarketSecurityPrediction = moneyMarketSecurityFactory
.getMoneyMarketSecurity();
moneyMarketSecurityMap.insert(moneyMarketSecurityPrediction);
moneyMarketSecurityPrediction.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
Listing listingPrediction = new Listing();
listingPrediction = listingBuilder.byUpdateMms(moneyMarketSecurityPrediction, listingPrediction);
listingImdg.insert(listingPrediction);
listingPrediction.setWorkflowStatus(ru.spcex.platform.enumeration.Status.Blocked.getKey());
CommonDeleteRequest commonDeleteRequest = new CommonDeleteRequest();
commonDeleteRequest.setId(ID);
//ACT
String jsonString = getJsonStringForDELETE(commonDeleteRequest, ID);
addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName()));
moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId());
MONEY_MARKET_SECURITY_MATCHER.assertMatch(moneyMarketSecurityResult, moneyMarketSecurityPrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", moneyMarketSecurityPrediction.getId()));
listingPrediction.setId(listingResult.getId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
/**
* {@link MoneyMarketSecurityService#updateMoneyMarket(BaseRequest)}
* Тест проверяет обновление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.<br>
* Входной запрос {@link MoneyMarketSecurityUpdateRequest}:<br>
* {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg<br>
* {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB<br>
* {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0<br>
* {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)<br>
*/
@Test
public void testUpdateMoneyMarket() {
clearAllInImdg(moneyMarketSecurityMap);
final String TOPIC = Consts.DESTINATION_MONEY_MARKET_SECURITY_UPDATE;
String updateTermType = "updateTermType";
Double updateLotSize = 12.0;
final MoneyMarketSecurity moneyMarketSecurityPrediction = moneyMarketSecurityFactory
.getMoneyMarketSecurity();
moneyMarketSecurityMap.insert(moneyMarketSecurityPrediction);
moneyMarketSecurityPrediction.setTermType(updateTermType);
moneyMarketSecurityPrediction.setLotSize(BigDecimal.valueOf(updateLotSize));
final MoneyMarketSecurityUpdateRequest keyRequest = moneyMarketSecurityFactory
.getMoneyMarketSecurityUpdateRequest();
keyRequest.setTermType(updateTermType);
keyRequest.setLotSize(updateLotSize);
Listing listingPrediction = new Listing();
listingPrediction = listingBuilder.byUpdateMms(moneyMarketSecurityPrediction, listingPrediction);
listingImdg.insert(listingPrediction);
listingPrediction.setLotSize(BigDecimal.valueOf(updateLotSize));
//ACT
String jsonString = getJsonStringForUPDATE(keyRequest, ID);
addRecordToKafka((MockConsumer) moneyMarketSecurityService.getConsumer(), TOPIC, PARTITION, 0, jsonString);
//ASSERT
waitingWhenAddedRecordAndCheckIt(ID, mockProducer, producerRecord);
MoneyMarketSecurity moneyMarketSecurityResult = moneyMarketSecurityMap.getSingleObjectBySQL(String.format("fullName = %s", moneyMarketSecurityFactory.getFullName()));
moneyMarketSecurityPrediction.setId(moneyMarketSecurityResult.getId());
moneyMarketSecurityPrediction.setSecurityId(moneyMarketSecurityResult.getSecurityId());
MONEY_MARKET_SECURITY_MATCHER.assertMatch(moneyMarketSecurityResult, moneyMarketSecurityPrediction);
Listing listingResult = listingImdg.getSingleObjectByFieldValues(Map.of("securityId", moneyMarketSecurityPrediction.getId()));
listingPrediction.setId(listingResult.getId());
listingPrediction.setSecurityId(listingResult.getSecurityId());
LISTING_MATCHER.assertMatch(listingResult, listingPrediction);
}
}

View file

@ -0,0 +1,123 @@
package ru.spcex.clearing.securities.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.platform.messaging.service.Status.Success;
import static ru.spcex.clearing.securities.utils.MatcherFactory.usingIgnoringFieldsComparator;
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 clearAllInImdg(Imdg<T> imdg) {
Collection<T> values = imdg.getAllValues();
values.forEach(imdg::delete);
}
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;
}
}
}