send one package with all info about company to company-service

This commit is contained in:
etreshenkov 2023-06-06 17:39:53 +03:00
parent 5c2c600d46
commit 5129a05bb8
3 changed files with 105 additions and 180 deletions

View file

@ -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;
}

View file

@ -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<CompaniesRequestParam> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final KafkaSender kafkaSender;
private final Imdg<Company> 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<MemberCompany> 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;
}
}

View file

@ -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<CompaniesRequestParam> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final KafkaSender kafkaSender;
private final Imdg<Company> 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<MemberCompany> 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;
}
}