From 70cb5b3da905c2e324e3c971bc89e8c54970110b Mon Sep 17 00:00:00 2001 From: akulikov Date: Fri, 26 May 2023 19:18:12 +0300 Subject: [PATCH] pipeline for FOND and EQUITY complete --- .../gatewayapi/config/KafkaConfig.java | 63 ++++++ .../config/PipelineConfiguration.java | 7 +- .../logic/listings/CheckSecurityExist.java | 4 + .../logic/listings/PrepareCompanies.java | 5 +- .../listings/SendMessageToCompanyService.java | 119 ++++++++++ .../SendMessageToSecurityService.java | 208 ++++++++++++++++++ .../listings/ValidateIncomeSecurities.java | 74 ++++++- .../request/objects/IncomeSecurity.java | 18 +- 8 files changed, 483 insertions(+), 15 deletions(-) create mode 100644 clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/KafkaConfig.java create mode 100644 clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToCompanyService.java create mode 100644 clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToSecurityService.java diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/KafkaConfig.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/KafkaConfig.java new file mode 100644 index 000000000..343f9d5f6 --- /dev/null +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/KafkaConfig.java @@ -0,0 +1,63 @@ +package ru.spcex.clearing.gatewayapi.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 org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +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 KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean + public Consumer createConsumer(GatewayApiSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(GatewayApiSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + + @Bean + public ProducerFactory pf(GatewayApiSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } +} diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/PipelineConfiguration.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/PipelineConfiguration.java index d4d8325d6..2b7dd4458 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/PipelineConfiguration.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/PipelineConfiguration.java @@ -4,6 +4,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import ru.spcex.clearing.gatewayapi.logic.Stage; import ru.spcex.clearing.gatewayapi.logic.listings.*; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.ImdgProvider; import java.util.ArrayList; @@ -13,12 +14,16 @@ import java.util.List; public class PipelineConfiguration { @Bean - public List> listingsRequestPipeline(ImdgProvider imdgProvider) { + public List> listingsRequestPipeline(ImdgProvider imdgProvider, KafkaSender kafkaSender) { List> pipeline = new ArrayList<>(); pipeline.add(new PrepareSecurities()); pipeline.add(new PrepareCompanies()); pipeline.add(new ValidateIncomeSecurities()); + pipeline.add(new ValidateIssuerCompanies()); pipeline.add(new CheckSecurityExist(imdgProvider)); + pipeline.add(new CheckCompanyExist(imdgProvider)); + pipeline.add(new SendMessageToCompanyService(kafkaSender, imdgProvider)); + pipeline.add(new SendMessageToSecurityService(kafkaSender, imdgProvider)); return pipeline; } diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/CheckSecurityExist.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/CheckSecurityExist.java index 9317597db..e97af6923 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/CheckSecurityExist.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/CheckSecurityExist.java @@ -174,6 +174,10 @@ public class CheckSecurityExist extends Stage { nominal.setAccruedCoupon(couponSchedule.getAccruedCoupon()); nominal.setCouponNumber(couponSchedule.getCouponNumber()); } + if (nominal.getAccruedCoupon() == null || nominal.getCouponNumber() == null) { + log.warn("Can't find accrued_coupon OR coupon_number for nominal, security.UUID {}", incomeSecurity.getId()); + nominal.setInvalidData(true); + } } } diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/PrepareCompanies.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/PrepareCompanies.java index 64e79a4a8..8dce664d2 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/PrepareCompanies.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/PrepareCompanies.java @@ -36,14 +36,13 @@ public class PrepareCompanies extends Stage { boolean securityFoundForCompany = false; for (IncomeSecurity incomeSecurity : param.getSecurities()) { if (incomeSecurity.getIssuerId().equals(issuerCompanyId)) { - incomeSecurity.getIssuerCompanyList().add(issuerCompany); + incomeSecurity.setIssuerCompany(issuerCompany); securityFoundForCompany = true; break; } } if (!securityFoundForCompany) { - log.warn("For company.UUID {} not found security, company skipped", issuerCompanyId); - issuerCompany.setInvalidData(true); + log.warn("For company.UUID {} not found security", issuerCompanyId); } } diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToCompanyService.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToCompanyService.java new file mode 100644 index 000000000..21bdbb229 --- /dev/null +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToCompanyService.java @@ -0,0 +1,119 @@ +package ru.spcex.clearing.gatewayapi.logic.listings; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.clearing.classes.statics.data.company.Company; +import ru.spcex.clearing.gatewayapi.logic.ProcessResult; +import ru.spcex.clearing.gatewayapi.logic.Stage; +import ru.spcex.clearing.gatewayapi.request.objects.IssuerCompany; +import ru.spcex.clearing.gatewayapi.request.objects.IssuerCompanyInfo; +import ru.spcex.clearing.gatewayapi.request.objects.IssuerCompanySymbols; +import ru.spcex.clearing.gatewayapi.request.objects.IssuerContact; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.company.*; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.CompanySymbol; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.List; +import java.util.Map; + +public class SendMessageToCompanyService extends Stage { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final KafkaSender kafkaSender; + private final Imdg companyImdg; + + + public SendMessageToCompanyService(KafkaSender kafkaSender, ImdgProvider imdgProvider) { + this.kafkaSender = kafkaSender; + this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); + } + + @Override + public ProcessResult process(ListingsRequestParam param) { + List companyList = param.getIssuerCompanyList(); + for (IssuerCompany company : companyList) { + if (company.isInvalidData()) continue; + + CompanyNewRequest companyNewRequest = new CompanyNewRequest(); + if (company.isAlreadyExist()) companyNewRequest.setId(company.getMapId()); + companyNewRequest.setShortName(company.getShortName()); + companyNewRequest.setFullName(company.getFullName()); + companyNewRequest.setWorkflowStatus(company.getWorkflowStatus()); + companyNewRequest.setCompanySymbol(CompanySymbol.UUID.getKey()); + companyNewRequest.setCompanySymbolValue(company.getId().toString()); + + IssuerCompanyInfo companyInfo = company.getIssuerCompanyInfo(); + + if (company.isAlreadyExist()) { + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_UPDATE, companyNewRequest); + } else { + kafkaSender.sendToQueueWaitForAnswer(Consts.DESTINATION_COMPANY_NEW, companyNewRequest); + Company companyFromImdg = companyImdg.getSingleObjectByFieldValues( + Map.of( + "shortName", company.getShortName(), + "fullName", company.getFullName(), + "workflowStatus", company.getWorkflowStatus() + ) + ); + if (companyFromImdg == null) { + log.warn("Can't insert company (UUID {}), skipped", company.getId()); + company.setInvalidData(true); + continue; + } + company.setAlreadyExist(true); + company.setMapId(companyFromImdg.getId()); + } + Long companyId = company.getMapId(); + + CompanyInfoUpdateRequest companyInfoUpdateRequest = new CompanyInfoUpdateRequest(); + companyInfoUpdateRequest.setId(companyId); + companyInfoUpdateRequest.setCountryCode(companyInfo.getCountryCode()); + companyInfoUpdateRequest.setLegalKind(companyInfo.getLegalKind()); + companyInfoUpdateRequest.setOrganizationType(companyInfo.getOrganizationType()); + companyInfoUpdateRequest.setResidence(companyInfo.getResidence()); + companyInfoUpdateRequest.setShortNameEng(companyInfo.getShortNameEng()); + companyInfoUpdateRequest.setFullNameEng(companyInfo.getFullNameEng()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_UPDATE, companyInfoUpdateRequest); + + for (IssuerCompanySymbols companySymbols : company.getIssuerCompanySymbolsList()) { + if (companySymbols.isInvalidData()) continue; + if (!company.isAlreadyExist() || !companySymbols.isAlreadyExist()) { + CompanySymbolNewRequest companySymbolNewRequest = new CompanySymbolNewRequest(); + companySymbolNewRequest.setCompanyId(companyId); + companySymbolNewRequest.setCompanySymbol(companySymbols.getCompanySymbol()); + companySymbolNewRequest.setCompanySymbolValue(companySymbols.getCompanySymbolValue()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_SYMBOL_NEW, companySymbolNewRequest); + } else { + CompanySymbolUpdateRequest companySymbolUpdateRequest = new CompanySymbolUpdateRequest(); + companySymbolUpdateRequest.setId(companySymbols.getMapId()); + companySymbolUpdateRequest.setCompanyId(companyId); + companySymbolUpdateRequest.setCompanySymbolValue(companySymbols.getCompanySymbolValue()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_SYMBOL_UPDATE, companySymbolUpdateRequest); + } + } + + for (IssuerContact contact : company.getIssuerContactList()) { + if (contact.isInvalidData()) continue; + if (!company.isAlreadyExist() || !contact.isAlreadyExist()) { + ContactNewRequest contactNewRequest = new ContactNewRequest(); + contactNewRequest.setCompanyId(companyId); + contactNewRequest.setContactType(contact.getContactType()); + contactNewRequest.setContactValue(contact.getContactValue()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_CONTACT_NEW, contactNewRequest); + } else { + ContactUpdateRequest contactUpdateRequest = new ContactUpdateRequest(); + contactUpdateRequest.setId(contact.getMapId()); + contactUpdateRequest.setContactType(contact.getContactType()); + contactUpdateRequest.setContactValue(contact.getContactValue()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_CONTACT_UPDATE, contactUpdateRequest); + } + } + } + + return null; + } +} diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToSecurityService.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToSecurityService.java new file mode 100644 index 000000000..fdfde2c46 --- /dev/null +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/SendMessageToSecurityService.java @@ -0,0 +1,208 @@ +package ru.spcex.clearing.gatewayapi.logic.listings; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity; +import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity; +import ru.spcex.clearing.gatewayapi.logic.ProcessResult; +import ru.spcex.clearing.gatewayapi.logic.Stage; +import ru.spcex.clearing.gatewayapi.request.objects.CouponSchedule; +import ru.spcex.clearing.gatewayapi.request.objects.IncomeListing; +import ru.spcex.clearing.gatewayapi.request.objects.IncomeSecurity; +import ru.spcex.clearing.gatewayapi.request.objects.Nominal; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.securitites.*; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.InstrumentType; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.List; +import java.util.Map; + +public class SendMessageToSecurityService extends Stage { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final KafkaSender kafkaSender; + private final Imdg fixedIncomeSecurityImdg; + private final Imdg equitySecurityImdg; + + + public SendMessageToSecurityService(KafkaSender kafkaSender, ImdgProvider imdgProvider) { + this.kafkaSender = kafkaSender; + this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class); + this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class); + } + + @Override + public ProcessResult process(ListingsRequestParam param) { + List securityList = param.getSecurities(); + for (IncomeSecurity security : securityList) { + if (security.isInvalidData()) continue; + if (security.getIssuerCompany() == null || security.getIssuerCompany().getMapId() == null) { + log.warn("Can't insert security.UUID {}, company not found for parameter issuerId", security.getId()); + security.setInvalidData(true); + continue; + } + Long securityId; + if (InstrumentType.BOND.equalsByKey(security.getInstrumentType())) { + if (!security.isAlreadyExist()) { + FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest = new FixedIncomeSecurityNewRequest(); + fixedIncomeSecurityNewRequest.setSecuritySymbol(security.getSecuritySymbol()); + fixedIncomeSecurityNewRequest.setShortName(security.getShortName()); + fixedIncomeSecurityNewRequest.setFullName(security.getFullName()); + fixedIncomeSecurityNewRequest.setIsin(security.getIsin()); + fixedIncomeSecurityNewRequest.setBondType(security.getBondType()); + fixedIncomeSecurityNewRequest.setNominalValue(security.getNominalValue()); + fixedIncomeSecurityNewRequest.setNominalCurrency(security.getNominalCurrency()); + fixedIncomeSecurityNewRequest.setMaturityDate(security.getMaturityDate()); + fixedIncomeSecurityNewRequest.setCouponFrequency(security.getCouponFrequency()); + fixedIncomeSecurityNewRequest.setIssuerId(security.getIssuerCompany().getMapId()); + fixedIncomeSecurityNewRequest.setShortNameEng(security.getShortNameEng()); + fixedIncomeSecurityNewRequest.setFullNameEng(security.getFullNameEng()); + fixedIncomeSecurityNewRequest.setWorkflowStatus(security.getWorkflowStatus()); + fixedIncomeSecurityNewRequest.setInstrumentType(security.getInstrumentType()); + kafkaSender.sendToQueueWaitForAnswer(Consts.DESTINATION_FIXED_INCOME_SECURITY_NEW, fixedIncomeSecurityNewRequest); + } else { + FixedIncomeSecurityUpdateRequest fixedIncomeSecurityUpdateRequest = new FixedIncomeSecurityUpdateRequest(); + fixedIncomeSecurityUpdateRequest.setId(security.getMapId()); + fixedIncomeSecurityUpdateRequest.setSecuritySymbol(security.getSecuritySymbol()); + fixedIncomeSecurityUpdateRequest.setShortName(security.getShortName()); + fixedIncomeSecurityUpdateRequest.setFullName(security.getFullName()); + fixedIncomeSecurityUpdateRequest.setIsin(security.getIsin()); + fixedIncomeSecurityUpdateRequest.setBondType(security.getBondType()); + fixedIncomeSecurityUpdateRequest.setNominalValue(security.getNominalValue()); + fixedIncomeSecurityUpdateRequest.setNominalCurrency(security.getNominalCurrency()); + fixedIncomeSecurityUpdateRequest.setMaturityDate(security.getMaturityDate()); + fixedIncomeSecurityUpdateRequest.setCouponFrequency(security.getCouponFrequency()); + fixedIncomeSecurityUpdateRequest.setIssuerId(security.getIssuerCompany().getMapId()); + fixedIncomeSecurityUpdateRequest.setShortNameEng(security.getShortNameEng()); + fixedIncomeSecurityUpdateRequest.setFullNameEng(security.getFullNameEng()); + fixedIncomeSecurityUpdateRequest.setWorkflowStatus(security.getWorkflowStatus()); + fixedIncomeSecurityUpdateRequest.setInstrumentType(security.getInstrumentType()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_SECURITY_UPDATE, fixedIncomeSecurityUpdateRequest); + } + } else { + if (security.isAlreadyExist()) { + EquitySecurityNewRequest equitySecurityNewRequest = new EquitySecurityNewRequest(); + equitySecurityNewRequest.setSecuritySymbol(security.getSecuritySymbol()); + equitySecurityNewRequest.setShortName(security.getShortName()); + equitySecurityNewRequest.setFullName(security.getFullName()); + equitySecurityNewRequest.setIsin(security.getIsin()); + equitySecurityNewRequest.setShareType(security.getShareType()); + equitySecurityNewRequest.setIssuerId(security.getIssuerCompany().getMapId()); + equitySecurityNewRequest.setShortNameEng(security.getShortNameEng()); + equitySecurityNewRequest.setFullNameEng(security.getFullNameEng()); + equitySecurityNewRequest.setWorkflowStatus(security.getWorkflowStatus()); + equitySecurityNewRequest.setInstrumentType(security.getInstrumentType()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_SECURITY_NEW, equitySecurityNewRequest); + } else { + EquitySecurityUpdateRequest equitySecurityUpdateRequest = new EquitySecurityUpdateRequest(); + equitySecurityUpdateRequest.setId(security.getMapId()); + equitySecurityUpdateRequest.setSecuritySymbol(security.getSecuritySymbol()); + equitySecurityUpdateRequest.setShortName(security.getShortName()); + equitySecurityUpdateRequest.setFullName(security.getFullName()); + equitySecurityUpdateRequest.setIsin(security.getIsin()); + equitySecurityUpdateRequest.setShareType(security.getShareType()); + equitySecurityUpdateRequest.setIssuerId(security.getIssuerCompany().getMapId()); + equitySecurityUpdateRequest.setShortNameEng(security.getShortNameEng()); + equitySecurityUpdateRequest.setFullNameEng(security.getFullNameEng()); + equitySecurityUpdateRequest.setWorkflowStatus(security.getWorkflowStatus()); + equitySecurityUpdateRequest.setInstrumentType(security.getInstrumentType()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_SECURITY_UPDATE, equitySecurityUpdateRequest); + } + } + + if (!security.isAlreadyExist()) { + Map predicate = Map.of( + "instrumentType", security.getInstrumentType(), + "securitySymbol", security.getSecuritySymbol(), + "bondType", security.getBondType(), + "workflowStatus", security.getWorkflowStatus(), + "shortName", security.getShortName() + ); + if (InstrumentType.BOND.equalsByKey(security.getInstrumentType())) { + FixedIncomeSecurity fixedIncomeSecurity = fixedIncomeSecurityImdg.getSingleObjectByFieldValues(predicate); + if (fixedIncomeSecurity != null) security.setMapId(fixedIncomeSecurity.getId()); + } else { + EquitySecurity equitySecurity = equitySecurityImdg.getSingleObjectByFieldValues(predicate); + if (equitySecurity != null) security.setMapId(equitySecurity.getId()); + } + if (security.getMapId() == null) { + log.warn("Can't insert security (UUID {}), skipped", security.getId()); + security.setInvalidData(true); + continue; + } + } + securityId = security.getMapId(); + + for (IncomeListing listing : security.getListingList()) { + if (listing.isInvalidData()) continue; + if (!security.isAlreadyExist() || !listing.isAlreadyExist()) { + ListingNewRequest listingNewRequest = new ListingNewRequest(); + listingNewRequest.setSecurityId(securityId); + listingNewRequest.setMarket(listing.getCode()); + listingNewRequest.setLotSize(listing.getLotSize()); + listingNewRequest.setTradingCurrency(listing.getTradingCurrency()); + listingNewRequest.setWorkflowStatus(listing.getWorkflowStatus()); + kafkaSender.sendRequestToQueue(Consts.LISTING_NEW, listingNewRequest); + } else { + ListingUpdateRequest listingUpdateRequest = new ListingUpdateRequest(); + listingUpdateRequest.setId(listing.getMapId()); + listingUpdateRequest.setSecurityId(securityId); + listingUpdateRequest.setMarket(listing.getCode()); + listingUpdateRequest.setLotSize(listing.getLotSize()); + listingUpdateRequest.setTradingCurrency(listing.getTradingCurrency()); + listingUpdateRequest.setWorkflowStatus(listing.getWorkflowStatus()); + kafkaSender.sendRequestToQueue(Consts.LISTING_UPDATE, listingUpdateRequest); + } + } + + for (CouponSchedule couponSchedule : security.getCouponScheduleList()) { + if (couponSchedule.isInvalidData()) continue; + if (!security.isAlreadyExist() || !couponSchedule.isAlreadyExist()) { + CouponPeriodNewRequest couponPeriodNewRequest = new CouponPeriodNewRequest(); + couponPeriodNewRequest.setSecurityId(securityId); + couponPeriodNewRequest.setCouponRate(couponSchedule.getCouponRate()); + couponPeriodNewRequest.setNumber(couponSchedule.getCouponNumber()); + couponPeriodNewRequest.setPeriodEndDate(couponSchedule.getPeriodEndDate()); + couponPeriodNewRequest.setPeriodStartDate(couponSchedule.getPeriodStartDate()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COUPON_PERIOD_NEW, couponPeriodNewRequest); + } else { + CouponPeriodUpdateRequest couponPeriodUpdateRequest = new CouponPeriodUpdateRequest(); + couponPeriodUpdateRequest.setId(couponSchedule.getMapId()); + couponPeriodUpdateRequest.setSecurityId(securityId); + couponPeriodUpdateRequest.setCouponRate(couponSchedule.getCouponRate()); + couponPeriodUpdateRequest.setNumber(couponSchedule.getCouponNumber()); + couponPeriodUpdateRequest.setPeriodEndDate(couponSchedule.getPeriodEndDate()); + couponPeriodUpdateRequest.setPeriodStartDate(couponSchedule.getPeriodStartDate()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COUPON_PERIOD_UPDATE, couponPeriodUpdateRequest); + } + } + + for (Nominal nominal : security.getNominalList()) { + if (nominal.isInvalidData()) continue; + if (!security.isAlreadyExist() || !nominal.isAlreadyExist()) { + FixedIncomeCashFlowNewRequest fixedIncomeCashFlowNewRequest = new FixedIncomeCashFlowNewRequest(); + fixedIncomeCashFlowNewRequest.setSecuritySymbol(security.getSecuritySymbol()); + fixedIncomeCashFlowNewRequest.setAccruedCoupon(nominal.getAccruedCoupon()); + fixedIncomeCashFlowNewRequest.setNominalValue(nominal.getNominal()); + fixedIncomeCashFlowNewRequest.setNumber(nominal.getCouponNumber()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_CASH_FLOW_NEW, fixedIncomeCashFlowNewRequest); + } else { + FixedIncomeCashFlowUpdateRequest fixedIncomeCashFlowUpdateRequest = new FixedIncomeCashFlowUpdateRequest(); + fixedIncomeCashFlowUpdateRequest.setId(nominal.getMapId()); + fixedIncomeCashFlowUpdateRequest.setSecuritySymbol(security.getSecuritySymbol()); + fixedIncomeCashFlowUpdateRequest.setAccruedCoupon(nominal.getAccruedCoupon()); + fixedIncomeCashFlowUpdateRequest.setNominalValue(nominal.getNominal()); + fixedIncomeCashFlowUpdateRequest.setNumber(nominal.getCouponNumber()); + kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_CASH_FLOW_UPDATE, fixedIncomeCashFlowUpdateRequest); + } + } + + } + + return null; + } +} diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/ValidateIncomeSecurities.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/ValidateIncomeSecurities.java index d884a39fd..3d8de36cf 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/ValidateIncomeSecurities.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/listings/ValidateIncomeSecurities.java @@ -5,8 +5,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import ru.spcex.clearing.gatewayapi.logic.ProcessResult; import ru.spcex.clearing.gatewayapi.logic.Stage; +import ru.spcex.clearing.gatewayapi.request.objects.CouponSchedule; import ru.spcex.clearing.gatewayapi.request.objects.IncomeListing; import ru.spcex.clearing.gatewayapi.request.objects.IncomeSecurity; +import ru.spcex.clearing.gatewayapi.request.objects.Nominal; import ru.spcex.platform.enumeration.BondType; import ru.spcex.platform.enumeration.InstrumentType; import ru.spcex.platform.enumeration.ShareType; @@ -26,6 +28,12 @@ public class ValidateIncomeSecurities extends Stage { for (IncomeSecurity incomeSecurity : securities) { if (incomeSecurity.isInvalidData()) continue; + if (incomeSecurity.getIssuerCompany() == null) { + log.warn("For security.UUID {} not found company, skipped", incomeSecurity.getId()); + incomeSecurity.setInvalidData(true); + continue; + } + String instrumentTypeStr = incomeSecurity.getInstrumentType(); InstrumentType instrumentType = IEnumKey.getEnumByKey(InstrumentType.class, instrumentTypeStr); if (instrumentType != InstrumentType.BOND && instrumentType != InstrumentType.EQTY) { @@ -73,15 +81,77 @@ public class ValidateIncomeSecurities extends Stage { List listingList = incomeSecurity.getListingList(); for (IncomeListing listing : listingList) { + + if (StringUtils.isEmpty(listing.getCode())) { + log.warn("Empty or null code (market) for listing from security.UUID {}", incomeSecurity.getId()); + listing.setInvalidData(true); + } + + if (listing.getLotSize() == null) { + log.warn("lotSize is null for listing from security.UUID {}", incomeSecurity.getId()); + listing.setInvalidData(true); + } + + if (StringUtils.isEmpty(listing.getTradingCurrency())) { + log.warn("Empty or null tradingCurrency for listing from security.UUID {}", incomeSecurity.getId()); + listing.setInvalidData(true); + } + workflowStatusStr = listing.getWorkflowStatus(); workflowStatus = IEnumKey.getEnumByKey(WorkflowStatus.class, workflowStatusStr); if (workflowStatus == null) { log.warn("Invalid workflowStatus {} for listing from security.UUID {}", workflowStatusStr, incomeSecurity.getId()); listing.setInvalidData(true); - continue; } - // todo add validation after create new listing request + } + + incomeSecurity.setListingList( + listingList.stream() + .filter(incomeListing -> !incomeListing.isInvalidData()) + .collect(Collectors.toList()) + ); + + + for (CouponSchedule couponSchedule : incomeSecurity.getCouponScheduleList()) { + if ( + couponSchedule.getPeriodStartDate() == null || + couponSchedule.getPeriodEndDate() == null || + couponSchedule.getPeriodEndDate().isBefore(couponSchedule.getPeriodStartDate()) + ) { + log.warn("Invalid periodStartDate {} or periodEndDate {} for couponSchedule, security.UUID {}, skipped", + couponSchedule.getPeriodStartDate(), + couponSchedule.getPeriodEndDate(), + incomeSecurity.getId()); + couponSchedule.setInvalidData(true); + } + + if (couponSchedule.getCouponNumber() == null) { + log.warn("couponNumber is null for couponSchedule, security.UUID {}, skipped", incomeSecurity.getId()); + couponSchedule.setInvalidData(true); + } + + if (couponSchedule.getCouponRate() == null) { + log.warn("couponRate is null for couponSchedule, security.UUID {}, skipped", incomeSecurity.getId()); + couponSchedule.setInvalidData(true); + } + + if (couponSchedule.getAccruedCoupon() == null) { + log.warn("accruedCoupon is null for couponSchedule, security.UUID {}, skipped", incomeSecurity.getId()); + couponSchedule.setInvalidData(true); + } + } + + for (Nominal nominal : incomeSecurity.getNominalList()) { + if (nominal.getNominal() == null) { + log.warn("nominal (value) is null for nominal, security.UUID {}, skipped", incomeSecurity.getId()); + nominal.setInvalidData(true); + } + + if (nominal.getDate() == null) { + log.warn("date is null for nominal, security.UUID {}, skipped", incomeSecurity.getId()); + nominal.setInvalidData(true); + } } } diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/request/objects/IncomeSecurity.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/request/objects/IncomeSecurity.java index a68740573..8381148eb 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/request/objects/IncomeSecurity.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/request/objects/IncomeSecurity.java @@ -169,7 +169,7 @@ public class IncomeSecurity extends WithMapId { value = "Дополнительная информация\\Количество купонов в год", example = "1000" ) - private BigDecimal couponFrequency; + private Long couponFrequency; @@ -183,7 +183,8 @@ public class IncomeSecurity extends WithMapId { private List nominalList = new ArrayList<>(); @JsonIgnore - private List issuerCompanyList = new ArrayList<>(); + private IssuerCompany issuerCompany = null; + public UUID getId() { return id; @@ -321,11 +322,11 @@ public class IncomeSecurity extends WithMapId { this.couponType = couponType; } - public BigDecimal getCouponFrequency() { + public Long getCouponFrequency() { return couponFrequency; } - public void setCouponFrequency(BigDecimal couponFrequency) { + public void setCouponFrequency(Long couponFrequency) { this.couponFrequency = couponFrequency; } @@ -353,12 +354,11 @@ public class IncomeSecurity extends WithMapId { this.nominalList = nominalList; } - public List getIssuerCompanyList() { - return issuerCompanyList; + public IssuerCompany getIssuerCompany() { + return issuerCompany; } - public void setIssuerCompanyList(List issuerCompanyList) { - this.issuerCompanyList = issuerCompanyList; + public void setIssuerCompany(IssuerCompany issuerCompany) { + this.issuerCompany = issuerCompany; } - }