diff --git a/clearing-parent/securities-service/pom.xml b/clearing-parent/securities-service/pom.xml index 6d3e62da1..742105bf4 100644 --- a/clearing-parent/securities-service/pom.xml +++ b/clearing-parent/securities-service/pom.xml @@ -24,6 +24,10 @@ ru.spcex.platform platform-enum + + ru.spcex.clearing + clearing-utils + ru.spcex.platform platform-imdg-api-hazelcast-impl diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/component/ListingBuilder.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/component/ListingBuilder.java index 79d6f32d3..d8cec1b88 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/component/ListingBuilder.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/component/ListingBuilder.java @@ -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; + } } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java index 1669d0e82..34826f19b 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/KafkaConfig.java @@ -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 createConsumer(SecuritiesServiceSettings settings) { return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/ValidationConfig.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/ValidationConfig.java new file mode 100644 index 000000000..253f2395c --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/config/ValidationConfig.java @@ -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); + } + +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/EquitySecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/EquitySecurityService.java index 81af0e134..4744748b9 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/EquitySecurityService.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/EquitySecurityService.java @@ -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 kafkaQueue, - Producer kafkaProducer, - ImdgProvider imdgProvider, - ValidationProvider validation, - IMessageResolver messageResolver) { + Producer 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 userRequest){ + private RequestInfoUpdate newEquity(BaseRequest userRequest) { ImdgTransaction transaction = imdgProvider.newTransaction(); EquitySecurityNewRequest req = userRequest.getRequestPayload(); - Optional 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 userRequest){ + private RequestInfoUpdate updateEquity(BaseRequest 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 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; } } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/FixedIncomeSecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/FixedIncomeSecurityService.java new file mode 100644 index 000000000..bb5b8af08 --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/FixedIncomeSecurityService.java @@ -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 fixedIncomeSecurityImdg; + private final Imdg 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 kafkaQueue, + Producer 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 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 fixedIncomeImdg = transaction.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class); + Imdg 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 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 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; + } +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java index 51190d495..493124d92 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/service/cud/MoneyMarketSecurityService.java @@ -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 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 userRequest) { ImdgTransaction transaction = imdgProvider.newTransaction(); MoneyMarketSecurityNewRequest req = userRequest.getRequestPayload(); - Optional 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; } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java index 8f09e27c5..40371d3a4 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java @@ -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 mmsMap; private final Imdg equityMap; + private final Imdg 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 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 fixedIncomeNewValidator() { + return equityRequest -> { + ImdgValidationContext 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 fixedIncomeUpdateValidator() { + return equityRequest -> { + ImdgValidationContext 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 fixedIncomeDeleteValidator() { + return equityRequest -> { + ImdgValidationContext 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() ); }; diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/EquityNewValidationRule.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/EquityNewValidationRule.java index 31d2cb36a..6e61fb0c8 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/EquityNewValidationRule.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/EquityNewValidationRule.java @@ -14,7 +14,7 @@ public enum EquityNewValidationRule implements IValidationRule validate(ImdgValidationContext context) { EquitySecurityNewRequest action = context.getValidatedObject(); - if (StringUtils.hasText(action.getShortName())){ + if (!StringUtils.hasText(action.getShortName())){ return of(SecuritiesError.RequiredFieldIsEmpty, "shortName"); } if (action.getLotSize() == null) { diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/FixedIncomeNewValidationRule.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/FixedIncomeNewValidationRule.java new file mode 100644 index 000000000..fd08634c8 --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/FixedIncomeNewValidationRule.java @@ -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> { + RequiredFieldIsNotEmpty() { + @Override + public Optional validate(ImdgValidationContext 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(); + } +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentTypeById.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentTypeById.java new file mode 100644 index 000000000..21b30f3af --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentTypeById.java @@ -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> { + private final IEnumId errorEnum; + private final String predictableInstrumentType; + + public IsValidInstrumentTypeById(IEnumId errorEnum, String predictableInstrumentType) { + this.errorEnum = errorEnum; + this.predictableInstrumentType = predictableInstrumentType; + } + + @Override + public Optional validate(ImdgValidationContext 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(); + } + } +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentType.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentTypeFromRequest.java similarity index 83% rename from clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentType.java rename to clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentTypeFromRequest.java index 04c8878d7..68b8bb857 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentType.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/IsValidInstrumentTypeFromRequest.java @@ -9,11 +9,11 @@ import ru.spcex.platform.utils.validation.IValidationRule; import java.util.Optional; -public class IsValidInstrumentType implements IValidationRule> { +public class IsValidInstrumentTypeFromRequest implements IValidationRule> { 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; } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsNewValidationRule.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsNewValidationRule.java index de4eaff95..5c19b309c 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsNewValidationRule.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsNewValidationRule.java @@ -14,7 +14,7 @@ public enum MmsNewValidationRule implements IValidationRule validate(ImdgValidationContext 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 validate(ImdgValidationContext context) { WithSecuritySymbol action = context.getValidatedObject(); Imdg 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( diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/HazelcastInstanceTestConfiguration.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/HazelcastInstanceTestConfiguration.java deleted file mode 100644 index c8294909c..000000000 --- a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/HazelcastInstanceTestConfiguration.java +++ /dev/null @@ -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(); - } - -} \ No newline at end of file diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/HazelcastServiceTestConfiguration.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/ImdgTestConfig.java similarity index 62% rename from clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/HazelcastServiceTestConfiguration.java rename to clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/ImdgTestConfig.java index 3af3f0514..95f01963c 100644 --- a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/HazelcastServiceTestConfiguration.java +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/ImdgTestConfig.java @@ -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; } -} \ No newline at end of file +} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/KafkaTestConfig.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/KafkaTestConfig.java new file mode 100644 index 000000000..5f54dac64 --- /dev/null +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/config/KafkaTestConfig.java @@ -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 kafkaProducer, + ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .producer(kafkaProducer) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } + + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public MockConsumer createTestConsumer() { + return new MockConsumer<>(OffsetResetStrategy.EARLIEST); + } +} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/AbstractServiceTest.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/AbstractServiceTest.java new file mode 100644 index 000000000..391580ce8 --- /dev/null +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/AbstractServiceTest.java @@ -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_MATCHER = usingIgnoringFieldsComparator("created", "updated"); + protected static AtomicLong currentId = new AtomicLong(0L); + protected Imdg moneyMarketSecurityMap; + protected Imdg listingImdg; + protected Imdg equitySecurityImdg; + protected Imdg fixedIncomeSecurityImdg; + protected final ListingBuilder listingBuilder = new ListingBuilder(); + @Captor + protected ArgumentCaptor producerRecord; + @MockBean + protected MockProducer 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()); + } +} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/EquitySecurityServiceTest.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/EquitySecurityServiceTest.java new file mode 100644 index 000000000..960212afc --- /dev/null +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/EquitySecurityServiceTest.java @@ -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 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)}
+ * Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link MoneyMarketSecurityNewRequest}:
+ * {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0
+ * {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
+ * {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
+ * {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0
+ * {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L
+ * {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"
+ * {@link MoneyMarketSecurityNewRequest#fullName} - "estat"
+ */ + @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.
+ * Входной запрос {@link MoneyMarketSecurityUpdateRequest}:
+ * {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg
+ * {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe
+ * {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB
+ * {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0
+ * {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)
+ */ + @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)}
+ * Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link CommonDeleteRequest}:
+ * {@link CommonDeleteRequest#id} - Идентификатор записи
+ */ + @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; + } +} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/FixedIncomeSecurityServiceTest.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/FixedIncomeSecurityServiceTest.java new file mode 100644 index 000000000..9e19f531e --- /dev/null +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/FixedIncomeSecurityServiceTest.java @@ -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 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)}
+ * Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link MoneyMarketSecurityNewRequest}:
+ * {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0
+ * {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
+ * {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
+ * {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0
+ * {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L
+ * {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"
+ * {@link MoneyMarketSecurityNewRequest#fullName} - "estat"
+ */ + @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.
+ * Входной запрос {@link MoneyMarketSecurityUpdateRequest}:
+ * {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg
+ * {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe
+ * {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB
+ * {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0
+ * {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)
+ */ + @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)}
+ * Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link CommonDeleteRequest}:
+ * {@link CommonDeleteRequest#id} - Идентификатор записи
+ */ + @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; + } +} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityFactory.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityFactory.java index 86675c716..fd66d5950 100644 --- a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityFactory.java +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityFactory.java @@ -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; diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityServiceDisabled.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityServiceDisabled.java deleted file mode 100644 index 971c65732..000000000 --- a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityServiceDisabled.java +++ /dev/null @@ -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 kafkaMockQueue; - private final MockProducer 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 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)}
- * Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
- * Входной запрос {@link MoneyMarketSecurityNewRequest}:
- * {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0
- * {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
- * {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
- * {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0
- * {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L
- * {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"
- * {@link MoneyMarketSecurityNewRequest#fullName} - "estat"
- */ -// @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 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 iMap = hz.getMap(MAP); - - Object waiter = new Object(); - - String listenerID = iMap.addEntryListener((EntryAddedListener) entryEvent -> { - System.out.println("Checking Equality"); - MoneyMarketSecurity moneyMarketSecurityResult = entryEvent.getValue(); - moneyMarketSecurityPrediction.setId(entryEvent.getKey()); - - MatcherFactory.Matcher 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)}
- * Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
- * Входной запрос {@link CommonDeleteRequest}:
- * {@link CommonDeleteRequest#id} - Идентификатор записи
- */ -// @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 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 iMap = hz.getMap(MAP); - iMap.put(ID, initialMoneyMarketSecurity); - - Object waiter = new Object(); - String listenerID = iMap.addEntryListener((EntryRemovedListener) 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.
- * Входной запрос {@link MoneyMarketSecurityUpdateRequest}:
- * {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg
- * {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe
- * {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB
- * {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0
- * {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)
- */ -// @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 consumerRecord = - new ConsumerRecord<>(TOPIC, PARTITION, 0, "key", jsonRequest); - - kafkaMockQueue = new MockConsumer<>(OffsetResetStrategy.EARLIEST); - initKafka(TOPIC, PARTITION, OFFSET); - - kafkaMockQueue.schedulePollTask(() -> kafkaMockQueue.addRecord(consumerRecord)); - - IMap iMap = hz.getMap(MAP); - iMap.put(ID, initialMoneyMarketSecurity); - - isUsed = false; - - Object waiter = new Object(); - String listenerID = iMap.addEntryListener((EntryUpdatedListener) entryEvent -> { - System.out.println("Checking Equality.."); - MoneyMarketSecurity moneyMarketSecurityResult = entryEvent.getValue(); - - MatcherFactory.Matcher 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); - } - -} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityServiceTest.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityServiceTest.java new file mode 100644 index 000000000..2f95ecb3e --- /dev/null +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/service/MoneyMarketSecurityServiceTest.java @@ -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 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 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)}
+ * Тест проверяет генерацию сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link MoneyMarketSecurityNewRequest}:
+ * {@link MoneyMarketSecurityNewRequest#lotSize} - 1.0
+ * {@link MoneyMarketSecurityNewRequest#startDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
+ * {@link MoneyMarketSecurityNewRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z.
+ * {@link MoneyMarketSecurityNewRequest#nominalValue} - 1.0
+ * {@link MoneyMarketSecurityNewRequest#nominalCurrency} - 1L
+ * {@link MoneyMarketSecurityNewRequest#instrumentType} - "sdgaga"
+ * {@link MoneyMarketSecurityNewRequest#fullName} - "estat"
+ */ + @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)}
+ * Тест проверяет удаление сущности {@link MoneyMarketSecurity} в Hazelcast при передаче из Apache Kafka.
+ * Входной запрос {@link CommonDeleteRequest}:
+ * {@link CommonDeleteRequest#id} - Идентификатор записи
+ */ + @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.
+ * Входной запрос {@link MoneyMarketSecurityUpdateRequest}:
+ * {@link MoneyMarketSecurityUpdateRequest#fullName} - sfgsdfg
+ * {@link MoneyMarketSecurityUpdateRequest#instrumentType} - qwerqwe
+ * {@link MoneyMarketSecurityUpdateRequest#nominalCurrency} - RUB
+ * {@link MoneyMarketSecurityUpdateRequest#nominalValue} - 2.0
+ * {@link MoneyMarketSecurityUpdateRequest#endDate} - текущее колличество секунд с 1970-01-01T00:00:00Z + (24 * 3600)
+ */ + @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); + } + +} diff --git a/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/utils/TestUtils.java b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/utils/TestUtils.java new file mode 100644 index 000000000..c34e15531 --- /dev/null +++ b/clearing-parent/securities-service/src/test/java/ru/spcex/clearing/securities/utils/TestUtils.java @@ -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> BASE_REQUEST_MATCHER = usingIgnoringFieldsComparator(); + private static final ObjectMapper objectMapper = new ObjectMapper(); + + public static void waitingWhenAddedRecordAndCheckIt(Long id, MockProducer mockProducer, ArgumentCaptor producerRecord) { + BaseRequest 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 baseRequestResult = (BaseRequest) 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 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 String getJsonStringForNew(T accountRequest, long id) { + return getJsonBaseRequest(accountRequest, id, ActionType.NEW); + } + + public static String getJsonStringForUPDATE(T accountRequest, long id) { + return getJsonBaseRequest(accountRequest, id, ActionType.UPDATE); + } + + public static String getJsonStringForDELETE(T accountRequest, long id) { + return getJsonBaseRequest(accountRequest, id, ActionType.DELETE); + } + + private static String getJsonBaseRequest(T accountRequest, long id, ActionType actionType) { + BaseRequest 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 void clearAllInImdg(Imdg imdg) { + Collection values = imdg.getAllValues(); + values.forEach(imdg::delete); + } + + public static class FutureRecordMetadata implements Future { + @Override + public boolean cancel(boolean mayInterruptIfRunning) { + return false; + } + + @Override + public boolean isCancelled() { + return false; + } + + @Override + public boolean isDone() { + return false; + } + + @Override + public RecordMetadata get() throws InterruptedException, ExecutionException { + return null; + } + + @Override + public RecordMetadata get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { + return null; + } + } +}