From 5129a05bb820da4c41f867481ce322d187d354bb Mon Sep 17 00:00:00 2001 From: etreshenkov Date: Tue, 6 Jun 2023 17:39:53 +0300 Subject: [PATCH] send one package with all info about company to company-service --- .../config/ProcessorConfiguration.java | 2 +- .../companies/CompaniesKafkaMessenger.java | 104 ++++++++++ ...geToCompanyServiceWithMemberCompanies.java | 179 ------------------ 3 files changed, 105 insertions(+), 180 deletions(-) create mode 100644 clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/CompaniesKafkaMessenger.java delete mode 100644 clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/SendMessageToCompanyServiceWithMemberCompanies.java diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/ProcessorConfiguration.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/ProcessorConfiguration.java index 006799c35..e994af56a 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/ProcessorConfiguration.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/config/ProcessorConfiguration.java @@ -55,7 +55,7 @@ public class ProcessorConfiguration { pipeline.add(new PrepareMemberCompanies()); pipeline.add(new ValidateMemberCompanies()); pipeline.add(new CheckMemberCompanyExist(imdgProvider)); - pipeline.add(new SendMessageToCompanyServiceWithMemberCompanies(kafkaSender, imdgProvider)); + pipeline.add(new CompaniesKafkaMessenger(kafkaSender, imdgProvider)); processor.setPipeline(pipeline); return processor; } diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/CompaniesKafkaMessenger.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/CompaniesKafkaMessenger.java new file mode 100644 index 000000000..f00b3266b --- /dev/null +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/CompaniesKafkaMessenger.java @@ -0,0 +1,104 @@ +package ru.spcex.clearing.gatewayapi.logic.companies; + +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.*; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.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; + +/** + * Отправка сообщения к company-service на создание/обновление company + * todo Пока убрал совсем не рабочую логику с sendToQueueWaitForAnswer + */ +public class CompaniesKafkaMessenger extends Stage { + private final Logger log = LoggerFactory.getLogger(getClass()); + + private final KafkaSender kafkaSender; + private final Imdg companyImdg; + + public CompaniesKafkaMessenger(KafkaSender kafkaSender, ImdgProvider imdgProvider) { + this.kafkaSender = kafkaSender; + this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); + } + + @Override + public ProcessResult process(CompaniesRequestParam param) { + List companyList = param.getMemberCompanyList(); + for (MemberCompany company : companyList) { + if (company.isInvalidData()) continue; + + MultiCompanyRequest multiCompanyRequest = new MultiCompanyRequest(); + + CompanyNewRequest companyNewRequest = new CompanyNewRequest(); + companyNewRequest.setShortName(company.getShortName()); + companyNewRequest.setFullName(company.getFullName()); + companyNewRequest.setCompanySymbol(CompanySymbol.UUID.getKey()); + companyNewRequest.setCompanySymbolValue(company.getId().toString()); + multiCompanyRequest.setCompany(companyNewRequest); + + MemberCompanyInfo companyInfo = company.getCompanyInfo(); + CompanyInfoUpdateRequest companyInfoUpdateRequest = new CompanyInfoUpdateRequest(); + companyInfoUpdateRequest.setCountryCode("RUS"); // todo в текущей версии ТЗ присылается в цифровом обозначении + companyInfoUpdateRequest.setProfessionalSign(companyInfo.getProfessionalSign()); + companyInfoUpdateRequest.setLegalKind(companyInfo.getLegalKind()); + companyInfoUpdateRequest.setOrganizationType(companyInfo.getOrganizationType()); + companyInfoUpdateRequest.setResidence(companyInfo.getResidence()); + multiCompanyRequest.setCompanyInfo(companyInfoUpdateRequest); + + for (MemberCompanySymbols companySymbols : company.getMemberCompanySymbolsList()) { + if (companySymbols.isInvalidData()) continue; + CompanySymbolNewRequest companySymbolNewRequest = new CompanySymbolNewRequest(); + companySymbolNewRequest.setCompanySymbol(companySymbols.getCompanySymbol()); + companySymbolNewRequest.setCompanySymbolValue(companySymbols.getCompanySymbolValue()); + multiCompanyRequest.getCompanySymbols().add(companySymbolNewRequest); + } + + for (MemberContact contact : company.getContactList()) { + if (contact.isInvalidData()) continue; + ContactNewRequest contactNewRequest = new ContactNewRequest(); + contactNewRequest.setContactType(contact.getContactType()); + contactNewRequest.setContactValue(contact.getContactValue()); + multiCompanyRequest.getContacts().add(contactNewRequest); + } + + for (MemberProfileDocument profileDocument : company.getProfileDocumentList()) { + if (profileDocument.isInvalidData()) continue; + ProfileDocumentNewRequest profileDocumentNewRequest = new ProfileDocumentNewRequest(); + profileDocumentNewRequest.setDocumentType(profileDocument.getDocumentType()); + profileDocumentNewRequest.setIssueDate(profileDocument.getIssueDate()); + profileDocumentNewRequest.setIssuer(profileDocument.getIssuer()); + profileDocumentNewRequest.setNumber(profileDocument.getNumber()); + profileDocumentNewRequest.setValidToDate(profileDocument.getValidToDate()); + profileDocumentNewRequest.setLink(profileDocument.getLink()); + multiCompanyRequest.getProfileDocuments().add(profileDocumentNewRequest); + } + + for (MemberClient client : company.getClientList()) { + if (client.isInvalidData()) continue; + + if (!company.isAlreadyExist() || !client.isAlreadyExist()) { + ClientCodeNewRequest clientCodeNewRequest = new ClientCodeNewRequest(); + clientCodeNewRequest.setCode(client.getClientCode()); + clientCodeNewRequest.setMoneyAccountId(client.getMoneyAccountId()); + clientCodeNewRequest.setDepoAccountId(client.getDepoAccountId()); + multiCompanyRequest.getClientCodes().add(clientCodeNewRequest); + } + } + + kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_MULTIREQUEST, multiCompanyRequest); + } + + return null; + } +} diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/SendMessageToCompanyServiceWithMemberCompanies.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/SendMessageToCompanyServiceWithMemberCompanies.java deleted file mode 100644 index f095ea803..000000000 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/logic/companies/SendMessageToCompanyServiceWithMemberCompanies.java +++ /dev/null @@ -1,179 +0,0 @@ -package ru.spcex.clearing.gatewayapi.logic.companies; - -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.*; -import ru.spcex.clearing.imdg.IMDGDistributedNames; -import ru.spcex.clearing.platform.messaging.domain.Consts; -import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeNewRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.account.ClientCodeUpdateRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.company.*; -import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; -import ru.spcex.platform.enumeration.CompanySymbol; -import ru.spcex.platform.enumeration.ServiceStatus; -import ru.spcex.platform.enumeration.WorkflowStatus; -import ru.spcex.platform.imdg.api.Imdg; -import ru.spcex.platform.imdg.api.ImdgProvider; -import ru.spcex.platform.utils.enumeration.IEnumKey; - -import java.util.List; -import java.util.Map; - -/** - * Отправка сообщения к company-service на создание/обновление company - * todo Пока убрал совсем не рабочую логику с sendToQueueWaitForAnswer - */ -public class SendMessageToCompanyServiceWithMemberCompanies extends Stage { - private final Logger log = LoggerFactory.getLogger(getClass()); - - private final KafkaSender kafkaSender; - private final Imdg companyImdg; - - public SendMessageToCompanyServiceWithMemberCompanies(KafkaSender kafkaSender, ImdgProvider imdgProvider) { - this.kafkaSender = kafkaSender; - this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); - } - @Override - public ProcessResult process(CompaniesRequestParam param) { - List companyList = param.getMemberCompanyList(); - for (MemberCompany company : companyList) { - if (company.isInvalidData()) continue; - - String companyWorkflowStatus = WorkflowStatus.Active.getKey(); - if (!company.getCompanyClearingCategoryList().isEmpty()) { - MemberCompanyClearingCategory clearingCategory = company.getCompanyClearingCategoryList().get(0); - ServiceStatus categoryStatus = IEnumKey.getEnumByKey(ServiceStatus.class, clearingCategory.getWorkflowStatus()); - if (categoryStatus == ServiceStatus.Active || categoryStatus == ServiceStatus.Appl || categoryStatus == ServiceStatus.Reopened) { - companyWorkflowStatus = WorkflowStatus.Active.getKey(); - } else { - companyWorkflowStatus = WorkflowStatus.Blocked.getKey(); - } - } - - CompanyNewRequest companyNewRequest = new CompanyNewRequest(); - if (company.isAlreadyExist()) companyNewRequest.setId(company.getMapId()); - companyNewRequest.setShortName(company.getShortName()); - companyNewRequest.setFullName(company.getFullName()); - companyNewRequest.setCompanySymbol(CompanySymbol.UUID.getKey()); - companyNewRequest.setCompanySymbolValue(company.getId().toString()); - - MemberCompanyInfo companyInfo = company.getCompanyInfo(); - - if (company.isAlreadyExist()) { - kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_UPDATE, companyNewRequest); - } else { - kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_NEW, companyNewRequest); - Company companyFromImdg = companyImdg.getSingleObjectByFieldValues( - Map.of( - "shortName", company.getShortName(), - "fullName", company.getFullName(), - "workflowStatus", companyWorkflowStatus - ) - ); - 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("RUS"); // todo в текущей версии ТЗ присылается в цифровом обозначении - companyInfoUpdateRequest.setProfessionalSign(companyInfo.getProfessionalSign()); - companyInfoUpdateRequest.setLegalKind(companyInfo.getLegalKind()); - companyInfoUpdateRequest.setOrganizationType(companyInfo.getOrganizationType()); - companyInfoUpdateRequest.setResidence(companyInfo.getResidence()); - kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_UPDATE, companyInfoUpdateRequest); - - for (MemberCompanySymbols companySymbols : company.getMemberCompanySymbolsList()) { - 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 (MemberContact contact : company.getContactList()) { - 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); - } - } - - for (MemberProfileDocument profileDocument : company.getProfileDocumentList()) { - if (profileDocument.isInvalidData()) continue; - if (!company.isAlreadyExist() || !profileDocument.isAlreadyExist()) { - ProfileDocumentNewRequest profileDocumentNewRequest = new ProfileDocumentNewRequest(); - profileDocumentNewRequest.setCompanyId(companyId); - profileDocumentNewRequest.setDocumentType(profileDocument.getDocumentType()); - profileDocumentNewRequest.setIssueDate(profileDocument.getIssueDate()); - profileDocumentNewRequest.setIssuer(profileDocument.getIssuer()); - profileDocumentNewRequest.setNumber(profileDocument.getNumber()); - profileDocumentNewRequest.setValidToDate(profileDocument.getValidToDate()); - profileDocumentNewRequest.setLink(profileDocument.getLink()); - kafkaSender.sendRequestToQueue(Consts.DESTINATION_PROFILE_DOCUMENT_NEW, profileDocumentNewRequest); - } else { - ProfileDocumentUpdateRequest profileDocumentUpdateRequest = new ProfileDocumentUpdateRequest(); - profileDocumentUpdateRequest.setId(profileDocument.getMapId()); - profileDocumentUpdateRequest.setCompanyId(companyId); - profileDocumentUpdateRequest.setDocumentType(profileDocument.getDocumentType()); - profileDocumentUpdateRequest.setIssueDate(profileDocument.getIssueDate()); - profileDocumentUpdateRequest.setIssuer(profileDocument.getIssuer()); - profileDocumentUpdateRequest.setNumber(profileDocument.getNumber()); - profileDocumentUpdateRequest.setValidToDate(profileDocument.getValidToDate()); - profileDocumentUpdateRequest.setLink(profileDocument.getLink()); - kafkaSender.sendRequestToQueue(Consts.DESTINATION_PROFILE_DOCUMENT_NEW, profileDocumentUpdateRequest); - } - } - - for (MemberClient client : company.getClientList()) { - if (client.isInvalidData()) continue; - - if (!company.isAlreadyExist() || !client.isAlreadyExist()) { - ClientCodeNewRequest clientCodeNewRequest = new ClientCodeNewRequest(); - clientCodeNewRequest.setCompanyId(companyId); - clientCodeNewRequest.setCode(client.getClientCode()); - clientCodeNewRequest.setMoneyAccountId(client.getMoneyAccountId()); - clientCodeNewRequest.setDepoAccountId(client.getDepoAccountId()); - kafkaSender.sendRequestToQueue(Consts.DESTINATION_CLIENT_CODE_NEW, clientCodeNewRequest); - } else { - ClientCodeUpdateRequest clientCodeUpdateRequest = new ClientCodeUpdateRequest(); - clientCodeUpdateRequest.setId(client.getMapId()); - clientCodeUpdateRequest.setCompanyId(companyId); - clientCodeUpdateRequest.setCode(client.getClientCode()); - clientCodeUpdateRequest.setMoneyAccountId(client.getMoneyAccountId()); - clientCodeUpdateRequest.setDepoAccountId(client.getDepoAccountId()); - kafkaSender.sendRequestToQueue(Consts.DESTINATION_CLIENT_CODE_NEW, clientCodeUpdateRequest); - } - } - } - - return null; - } -}