some refactroing gateway

This commit is contained in:
etreschenkov 2023-06-16 17:23:19 +03:00
parent 2d2a89264d
commit c15f3fba9d
22 changed files with 279 additions and 866 deletions

View file

@ -98,7 +98,7 @@ public class MultiCompanyService
public void afterPropertiesSet() {
callback(CompanyGatewayRequest.class)
.setFunction(request -> requestHelper.requestFunction(this::processMultiRequest, request))
.forDestination(Consts.DESTINATION_COMPANY_MULTIREQUEST, callbacks::put);
.forDestination(Consts.DESTINATION_COMPANY_GATEWAY_REQUEST, callbacks::put);
init();
}

View file

@ -5,7 +5,6 @@ import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.spcex.clearing.gatewayapi.logic.Processor;
import ru.spcex.clearing.gatewayapi.logic.Stage;
import ru.spcex.clearing.gatewayapi.logic.listings_fond.*;
import ru.spcex.clearing.gatewayapi.logic.listings_mm.MMListingsRequestParam;
import ru.spcex.clearing.gatewayapi.logic.listings_mm.PrepareExchangeInstruments;
import ru.spcex.clearing.gatewayapi.logic.listings_mm.SendMessageToSecurityServiceWithMMSecurities;
@ -18,18 +17,6 @@ import java.util.List;
@Configuration
public class ProcessorConfiguration {
@Qualifier("fondListingsRequestProcessor")
@Bean
public Processor<FondListingsRequestParam> fondListingsRequestProcessor(ImdgProvider imdgProvider, KafkaSender kafkaSender) {
Processor<FondListingsRequestParam> processor = new Processor<>();
List<Stage<FondListingsRequestParam>> pipeline = new ArrayList<>();
pipeline.add(new PrepareIncomeSecurities());
pipeline.add(new PrepareIssuerCompanies());
pipeline.add(new SendMessageToCompanyServiceWithIssuerCompanies(kafkaSender, imdgProvider));
pipeline.add(new SendMessageToSecurityServiceWithFondSecurities(kafkaSender, imdgProvider));
processor.setPipeline(pipeline);
return processor;
}
@Qualifier("mmListingsRequestProcessor")
@Bean

View file

@ -19,9 +19,9 @@ import ru.spcex.clearing.gatewayapi.controller.response.CommonResponse;
import ru.spcex.clearing.gatewayapi.exception.GatewayException;
import ru.spcex.clearing.gatewayapi.logic.ProcessResult;
import ru.spcex.clearing.gatewayapi.logic.Processor;
import ru.spcex.clearing.gatewayapi.logic.listings_fond.FondListingsRequestParam;
import ru.spcex.clearing.gatewayapi.logic.listings_mm.MMListingsRequestParam;
import ru.spcex.clearing.gatewayapi.service.CompanyProcessor;
import ru.spcex.clearing.gatewayapi.service.ListingFondProcessor;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
@ -37,22 +37,23 @@ public class GatewayController {
private final Logger log = LoggerFactory.getLogger(getClass());
private final IMessageResolver messageResolver;
private final Processor<FondListingsRequestParam> fondListingsRequestProcessor;
private final Processor<MMListingsRequestParam> mmListingsRequestProcessor;
private final CompanyProcessor companiesProcessor;
private final ListingFondProcessor listingFondProcessor;
private final ExecutorService executor;
private List<String> validTypes = Arrays.asList("DAY_START", "ON_DEMAND");
public GatewayController(
IMessageResolver messageResolver,
@Qualifier("fondListingsRequestProcessor") Processor<FondListingsRequestParam> fondListingsRequestProcessor,
@Qualifier("mmListingsRequestProcessor") Processor<MMListingsRequestParam> mmListingsRequestProcessor,
CompanyProcessor companiesProcessor, @Qualifier("gatewayExecutor") ExecutorService executor) {
CompanyProcessor companiesProcessor,
ListingFondProcessor listingFondProcessor,
@Qualifier("gatewayExecutor") ExecutorService executor) {
this.messageResolver = messageResolver;
this.fondListingsRequestProcessor = fondListingsRequestProcessor;
this.mmListingsRequestProcessor = mmListingsRequestProcessor;
this.companiesProcessor = companiesProcessor;
this.listingFondProcessor = listingFondProcessor;
this.executor = executor;
}
@ -71,14 +72,7 @@ public class GatewayController {
@ResponseBody
public CommonResponse listingsFond(@RequestBody FondListingsRequest request) {
// request.validate(validTypes, List.of(Section.FOND));
FondListingsRequestParam requestParam = new FondListingsRequestParam(request);
executor.submit(() -> {
ProcessResult processResult = fondListingsRequestProcessor.process(requestParam);
if (processResult.getError() != null) {
log.warn("Exception while process listingFind: {}", processResult.getError());
}
});
executor.submit(() -> listingFondProcessor.process(request));
return createResponse(request, true);
}
@ -124,9 +118,7 @@ public class GatewayController {
@ResponseBody
public CommonResponse companies(@RequestBody CompaniesRequest request) {
// request.validate(validTypes, Collections.emptyList());
executor.submit(() -> {
companiesProcessor.process(request);
});
executor.submit(() -> companiesProcessor.process(request));
return createResponse(request, true);
}

View file

@ -0,0 +1,7 @@
package ru.spcex.clearing.gatewayapi.controller.request;
import java.util.UUID;
public interface WithIssuerId {
UUID getIssuerId();
}

View file

@ -5,6 +5,7 @@ import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import io.swagger.annotations.ApiModelProperty;
import ru.spcex.clearing.gatewayapi.config.deserializers.LocalDateDeserializer;
import ru.spcex.clearing.gatewayapi.controller.request.WithIssuerId;
import ru.spcex.clearing.gatewayapi.controller.request.WithMapId;
import java.math.BigDecimal;
@ -13,7 +14,7 @@ import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
public class FondSecurity extends WithMapId {
public class FondSecurity extends WithMapId implements WithIssuerId {
@JsonProperty("id")
@ApiModelProperty(
value = "id (uuid)",
@ -208,6 +209,7 @@ public class FondSecurity extends WithMapId {
this.instrumentType = instrumentType;
}
@Override
public UUID getIssuerId() {
return issuerId;
}

View file

@ -1,43 +0,0 @@
package ru.spcex.clearing.gatewayapi.logic.listings_fond;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondListingsRequest;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondSecurity;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompany;
import java.util.ArrayList;
import java.util.List;
public class FondListingsRequestParam {
private FondListingsRequest fondListingsRequest;
private List<FondSecurity> securities = new ArrayList<>();
private List<IssuerCompany> issuerCompanyList = new ArrayList<>();
public FondListingsRequestParam(FondListingsRequest fondListingsRequest) {
this.fondListingsRequest = fondListingsRequest;
}
public FondListingsRequest getRequest() {
return fondListingsRequest;
}
public List<FondSecurity> getSecurities() {
return securities;
}
public List<IssuerCompany> getIssuerCompanyList() {
return issuerCompanyList;
}
public void setListingsRequest(FondListingsRequest fondListingsRequest) {
this.fondListingsRequest = fondListingsRequest;
}
public void setSecurities(List<FondSecurity> securities) {
this.securities = securities;
}
public void setIssuerCompanyList(List<IssuerCompany> issuerCompanyList) {
this.issuerCompanyList = issuerCompanyList;
}
}

View file

@ -1,84 +0,0 @@
package ru.spcex.clearing.gatewayapi.logic.listings_fond;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.*;
import ru.spcex.clearing.gatewayapi.errors.GatewayError;
import ru.spcex.clearing.gatewayapi.logic.ProcessResult;
import ru.spcex.clearing.gatewayapi.logic.Stage;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.stream.Collectors;
public class PrepareIncomeSecurities extends Stage<FondListingsRequestParam> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public ProcessResult process(FondListingsRequestParam param) {
List<FondSecurity> securities = param.getRequest().getSecurities();
Map<UUID, FondSecurity> securityMap = new HashMap<>();
for (FondSecurity fondSecurity : securities) {
if (securityMap.containsKey(fondSecurity.getId())) {
log.warn("For security.id {} found > 1 element from request, use first", fondSecurity.getId());
fondSecurity.setInvalidData(true);
continue;
}
securityMap.put(fondSecurity.getId(), fondSecurity);
}
List<CouponSchedule> couponSchedules = param.getRequest().getCouponSchedules();
for (CouponSchedule couponSchedule : couponSchedules) {
UUID couponScheduleUUID = couponSchedule.getSecurityId();
FondSecurity fondSecurity = securityMap.get(couponScheduleUUID);
if (fondSecurity == null) {
log.warn("For coupon_schedule.security_id {} not found security, skipped", couponScheduleUUID);
continue;
}
fondSecurity.getCouponScheduleList().add(couponSchedule);
}
List<Currency> currencyList = param.getRequest().getCurrencies();
for (Currency currency : currencyList) {
UUID currencyUUID = currency.getId();
// todo сделать после выяснения
}
List<IncomeListing> listingList = param.getRequest().getListingList();
for (IncomeListing listing : listingList) {
UUID listingUUID = listing.getSecurityId();
FondSecurity fondSecurity = securityMap.get(listingUUID);
if (fondSecurity == null) {
log.warn("For listing.security_id {} not found security, skipped", listingUUID);
continue;
}
fondSecurity.getListingList().add(listing);
}
List<Nominal> nominalList = param.getRequest().getNominalList();
for (Nominal nominal : nominalList) {
UUID nominalUUID = nominal.getSecurityId();
FondSecurity fondSecurity = securityMap.get(nominalUUID);
if (fondSecurity == null) {
log.warn("For nominal.security_id {} not found security, skipped", nominalUUID);
continue;
}
fondSecurity.getNominalList().add(nominal);
}
List<FondSecurity> filteredIncomeSecurities = securities.stream()
.filter(incomeSecurity -> !incomeSecurity.isInvalidData())
.collect(Collectors.toList());
if (filteredIncomeSecurities.isEmpty()) {
return new ProcessResult(GatewayError.SecurityNotFound, securities.stream().map(fondSecurity -> fondSecurity.getId().toString()).collect(Collectors.joining(", ")));
}
param.setSecurities(filteredIncomeSecurities);
return null;
}
}

View file

@ -1,99 +0,0 @@
package ru.spcex.clearing.gatewayapi.logic.listings_fond;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.*;
import ru.spcex.clearing.gatewayapi.logic.ProcessResult;
import ru.spcex.clearing.gatewayapi.logic.Stage;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.stream.Collectors;
public class PrepareIssuerCompanies extends Stage<FondListingsRequestParam> {
private final Logger log = LoggerFactory.getLogger(getClass());
@Override
public ProcessResult process(FondListingsRequestParam param) {
List<IssuerCompany> issuerCompanyList = param.getRequest().getIssuerCompanyList();
Map<UUID, IssuerCompany> issuerCompanyMap = new HashMap<>();
for (IssuerCompany issuerCompany : issuerCompanyList) {
if (issuerCompanyMap.containsKey(issuerCompany.getId())) {
log.warn("For company.id {} found > 1 element from request, use first", issuerCompany.getId());
issuerCompany.setInvalidData(true);
continue;
}
issuerCompanyMap.put(issuerCompany.getId(), issuerCompany);
}
for (IssuerCompany issuerCompany : issuerCompanyList) {
if (issuerCompany.isInvalidData()) continue;
UUID issuerCompanyId = issuerCompany.getId();
boolean securityFoundForCompany = false;
for (FondSecurity fondSecurity : param.getSecurities()) {
if (fondSecurity.getIssuerId().equals(issuerCompanyId)) {
fondSecurity.setIssuerCompany(issuerCompany);
securityFoundForCompany = true;
break;
}
}
if (!securityFoundForCompany) {
log.warn("For company.id {} not found security", issuerCompanyId);
}
}
List<IssuerCompanyInfo> issuerCompanyInfoList = param.getRequest().getIssuerCompanyInfoList();
for (IssuerCompanyInfo issuerCompanyInfo : issuerCompanyInfoList) {
UUID issuerCompanyInfoUUID = issuerCompanyInfo.getCompanyId();
IssuerCompany issuerCompany = issuerCompanyMap.get(issuerCompanyInfoUUID);
if (issuerCompany == null) {
log.warn("For company_info.company_id {} not found company, company_info skipped", issuerCompanyInfoUUID);
continue;
}
if (issuerCompany.getIssuerCompanyInfo() != null) {
log.warn("For company.id {} found > 1 company_info, use first", issuerCompanyInfoUUID);
continue;
}
issuerCompany.setIssuerCompanyInfo(issuerCompanyInfo);
}
List<IssuerCompanySymbols> issuerCompanySymbolsList = param.getRequest().getIssuerCompanySymbolsList();
for (IssuerCompanySymbols issuerCompanySymbols : issuerCompanySymbolsList) {
UUID issuerCompanySymbolsUUID = issuerCompanySymbols.getCompanyId();
IssuerCompany issuerCompany = issuerCompanyMap.get(issuerCompanySymbolsUUID);
if (issuerCompany == null) {
log.warn("For company_symbols ({} = {}) not found company, company_symbols skipped",
issuerCompanySymbols.getCompanySymbol(),
issuerCompanySymbols.getCompanySymbolValue());
issuerCompanySymbols.setInvalidData(true);
continue;
}
issuerCompany.getIssuerCompanySymbolsList().add(issuerCompanySymbols);
}
List<IssuerContact> issuerContactList = param.getRequest().getIssuerContactList();
for (IssuerContact issuerContact : issuerContactList) {
UUID issuerContactUUID = issuerContact.getCompanyId();
IssuerCompany issuerCompany = issuerCompanyMap.get(issuerContactUUID);
if (issuerCompany == null) {
log.warn("For contact ({} = {}) not found company, contact skipped",
issuerContact.getContactType(),
issuerContact.getContactValue());
issuerContact.setInvalidData(true);
continue;
}
issuerCompany.getIssuerContactList().add(issuerContact);
}
List<IssuerCompany> filteredIssuerCompany = issuerCompanyList.stream()
.filter(issuerCompany -> !issuerCompany.isInvalidData())
.collect(Collectors.toList());
param.setIssuerCompanyList(filteredIssuerCompany);
return null;
}
}

View file

@ -1,125 +0,0 @@
package ru.spcex.clearing.gatewayapi.logic.listings_fond;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.clearing.classes.statics.data.company.Company;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompany;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompanyInfo;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompanySymbols;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerContact;
import ru.spcex.clearing.gatewayapi.logic.ProcessResult;
import ru.spcex.clearing.gatewayapi.logic.Stage;
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;
/**
* Отправка сообщения к company-service на создание/обновление company
* todo Пока убрал совсем не рабочую логику с sendToQueueWaitForAnswer
*/
public class SendMessageToCompanyServiceWithIssuerCompanies extends Stage<FondListingsRequestParam> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final KafkaSender kafkaSender;
private final Imdg<Company> companyImdg;
public SendMessageToCompanyServiceWithIssuerCompanies(KafkaSender kafkaSender, ImdgProvider imdgProvider) {
this.kafkaSender = kafkaSender;
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
}
@Override
public ProcessResult process(FondListingsRequestParam param) {
List<IssuerCompany> 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.sendRequestToQueue(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 (id {}), 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);
if (companyInfo != null) {
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_INFO_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;
}
}

View file

@ -1,244 +0,0 @@
package ru.spcex.clearing.gatewayapi.logic.listings_fond;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.CouponSchedule;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondSecurity;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IncomeListing;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.Nominal;
import ru.spcex.clearing.gatewayapi.logic.ProcessResult;
import ru.spcex.clearing.gatewayapi.logic.Stage;
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.CompanySymbol;
import ru.spcex.platform.enumeration.InstrumentType;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.UUID;
/**
* Отправка сообщения к security-service на создание/обновление security
* todo Пока убрал совсем не рабочую логику с sendToQueueWaitForAnswer
*/
public class SendMessageToSecurityServiceWithFondSecurities extends Stage<FondListingsRequestParam> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final KafkaSender kafkaSender;
private final Imdg<FixedIncomeSecurity> fixedIncomeSecurityImdg;
private final Imdg<EquitySecurity> equitySecurityImdg;
private final Imdg<CompanySymbols> companySymbolsImdg;
private final Imdg<Company> companyImdg;
public SendMessageToSecurityServiceWithFondSecurities(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);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.companySymbolsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
}
@Override
public ProcessResult process(FondListingsRequestParam param) {
List<FondSecurity> securityList = param.getSecurities();
for (FondSecurity security : securityList) {
if (security.isInvalidData()) continue;
Long companyId;
if (security.getIssuerCompany() == null || security.getIssuerCompany().getMapId() == null) {
// Пытаемся найти компанию по UUID
UUID companyUUID = security.getIssuerId();
Collection<CompanySymbols> companySymbols = companySymbolsImdg.getCollectionObjectsByFieldValues(
Map.of(
"companySymbol", CompanySymbol.UUID.getKey(),
"companySymbolValue", companyUUID.toString()
)
);
if (companySymbols.isEmpty()) {
log.warn("Can't insert security.id {}, company not found for parameter issuerId", security.getId());
security.setInvalidData(true);
continue;
} else if (companySymbols.size() > 1) {
log.warn("For {} = {} found > 1 company_symbol, use first", CompanySymbol.UUID.getKey(), companyUUID);
}
// Проверяем, что в базе помимо реквизита есть компания
companyId = companySymbols.iterator().next().getCompanyId();
Company company = companyImdg.getSingleObjectByID(companyId);
if (company == null) {
log.warn("Can't insert security: not found company with id {}", companyId);
security.setInvalidData(true);
continue;
}
} else {
companyId = security.getIssuerCompany().getMapId();
}
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().longValue());
fixedIncomeSecurityNewRequest.setIssuerId(companyId);
fixedIncomeSecurityNewRequest.setShortNameEng(security.getShortNameEng());
fixedIncomeSecurityNewRequest.setFullNameEng(security.getFullNameEng());
fixedIncomeSecurityNewRequest.setWorkflowStatus(security.getWorkflowStatus());
fixedIncomeSecurityNewRequest.setInstrumentType(security.getInstrumentType());
kafkaSender.sendRequestToQueue(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().longValue());
fixedIncomeSecurityUpdateRequest.setIssuerId(companyId);
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(companyId);
equitySecurityNewRequest.setShortNameEng(security.getShortNameEng());
equitySecurityNewRequest.setFullNameEng(security.getFullNameEng());
equitySecurityNewRequest.setWorkflowStatus(security.getWorkflowStatus());
equitySecurityNewRequest.setInstrumentType(security.getInstrumentType());
kafkaSender.sendRequestToQueue(Consts.DESTINATION_EQUITY_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(companyId);
equitySecurityUpdateRequest.setShortNameEng(security.getShortNameEng());
equitySecurityUpdateRequest.setFullNameEng(security.getFullNameEng());
equitySecurityUpdateRequest.setWorkflowStatus(security.getWorkflowStatus());
equitySecurityUpdateRequest.setInstrumentType(security.getInstrumentType());
kafkaSender.sendRequestToQueue(Consts.DESTINATION_EQUITY_SECURITY_UPDATE, equitySecurityUpdateRequest);
}
}
if (!security.isAlreadyExist()) {
Map<String, String> predicate = Map.of(
"instrumentType", security.getInstrumentType(),
"securitySymbol", security.getSecuritySymbol(),
"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 (id {}), 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;
}
}

View file

@ -63,7 +63,7 @@ public class CompanyProcessor {
companyGatewayRequest.setProfileDocuments(profileDocumentNewRequests);
companyGatewayRequest.setClientCodes(clientCodeNewRequests);
kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_MULTIREQUEST, companyGatewayRequest);
kafkaSender.sendRequestToQueue(Consts.DESTINATION_COMPANY_GATEWAY_REQUEST, companyGatewayRequest);
}
}

View file

@ -0,0 +1,65 @@
package ru.spcex.clearing.gatewayapi.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.gatewayapi.controller.request.WithCompanyId;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondListingsRequest;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompany;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompanySymbols;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerContact;
import ru.spcex.clearing.gatewayapi.service.adapter.SecurityMkrRequestAdapter;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.SecurityMkrGatewayRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import java.util.List;
import java.util.UUID;
import java.util.stream.Collectors;
@Service
public class IssueCompanyService {
private final Logger log = LoggerFactory.getLogger(getClass());
private final SecurityMkrRequestAdapter securityMkrRequestAdapter;
private final KafkaSender kafkaSender;
public IssueCompanyService(SecurityMkrRequestAdapter securityMkrRequestAdapter,
KafkaSender kafkaSender) {
this.securityMkrRequestAdapter = securityMkrRequestAdapter;
this.kafkaSender = kafkaSender;
}
public void sendRequest(FondListingsRequest request) {
List<IssuerCompany> issuerCompanyList = request.getIssuerCompanyList();
for (IssuerCompany issuerCompany : issuerCompanyList) {
UUID companyId = issuerCompany.getId();
log.debug("Grouping by companyId: {}", companyId);
List<IssuerCompanySymbols> issuerCompanySymbolsList = groupByCompanyId(request.getIssuerCompanySymbolsList(), companyId);
List<IssuerContact> issuerContactList = groupByCompanyId(request.getIssuerContactList(), companyId);
List<CompanySymbolNewRequest> companySymbolNewRequests = issuerCompanySymbolsList.stream().
map(securityMkrRequestAdapter::toCompanySymbolRequest).toList();
List<ContactNewRequest> contactNewRequests = issuerContactList.stream().
map(securityMkrRequestAdapter::toContactRequest).toList();
SecurityMkrGatewayRequest securityMkrGatewayRequest = new SecurityMkrGatewayRequest();
securityMkrGatewayRequest.setCompanyNewRequest(securityMkrRequestAdapter.toCompanyRequest(issuerCompany));
securityMkrGatewayRequest.setCompanySymbolNewRequests(companySymbolNewRequests);
securityMkrGatewayRequest.setContactNewRequests(contactNewRequests);
request.getIssuerCompanyInfoList()
.stream()
.filter(info -> info.getCompanyId().equals(companyId)).findFirst()
.map(securityMkrRequestAdapter::toCompanyInfoRequest)
.ifPresent(securityMkrGatewayRequest::setCompanyInfoUpdateRequest);
kafkaSender.sendRequestToQueue(Consts.DESTINATION_ISSUER_COMPANY_GATEWAY_REQUEST, securityMkrGatewayRequest);
}
}
private <T extends WithCompanyId> List<T> groupByCompanyId(List<T> listToProcess, UUID companyId) {
return listToProcess.stream().filter(t -> t.getCompanyId().equals(companyId)).collect(Collectors.toList());
}
}

View file

@ -1,69 +1,20 @@
package ru.spcex.clearing.gatewayapi.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.gatewayapi.controller.request.WithSecurityId;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.*;
import ru.spcex.clearing.gatewayapi.service.adapter.SecurityRequestAdapter;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.EquitySecurityGatewayRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.FixedIncomeGatewayRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.SecurityGatewayRequest;
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 java.util.List;
import java.util.UUID;
import java.util.stream.Collectors;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondListingsRequest;
@Service
public class ListingFondProcessor {
private final Logger log = LoggerFactory.getLogger(getClass());
private final IssueCompanyService issueCompanyProcessor;
private final SecurityService securityService;
private final SecurityRequestAdapter securityRequestAdapter;
private final KafkaSender kafkaSender;
public ListingFondProcessor(SecurityRequestAdapter securityRequestAdapter,
KafkaSender kafkaSender) {
this.securityRequestAdapter = securityRequestAdapter;
this.kafkaSender = kafkaSender;
public ListingFondProcessor(IssueCompanyService issueCompanyProcessor, SecurityService securityService) {
this.issueCompanyProcessor = issueCompanyProcessor;
this.securityService = securityService;
}
public void process(FondListingsRequest request) {
List<FondSecurity> securityList = request.getSecurities();
for (FondSecurity security : securityList) {
UUID securityId = security.getId();
log.debug("Grouping by securityId: {}", securityId);
List<IncomeListing> incomeListings = groupByCompanyId(request.getListingList(), securityId);
List<CouponSchedule> couponSchedules = groupByCompanyId(request.getCouponSchedules(), securityId);
List<Nominal> nominals = groupByCompanyId(request.getNominalList(), securityId);
List<ListingNewRequest> listingRequests = incomeListings.stream().
map(securityRequestAdapter::toIncomeListing).toList();
List<CouponPeriodNewRequest> couponPeriodRequests = couponSchedules.stream().
map(securityRequestAdapter::toCouponPeriodRequest).toList();
List<FixedIncomeCashFlowNewRequest> cashFlowRequests = nominals.stream().
map(securityRequestAdapter::toFixedIncomeCashFlowRequest).toList();
SecurityGatewayRequest<?> securityGatewayRequest;
if (InstrumentType.BOND.equalsByKey(security.getInstrumentType())) {
FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest = securityRequestAdapter.toFixedIncomeSecurityRequest(security);
securityGatewayRequest = new FixedIncomeGatewayRequest(fixedIncomeSecurityNewRequest);
} else {
EquitySecurityNewRequest equitySecurityNewRequest = securityRequestAdapter.toEquitySecurityNewRequest(security);
securityGatewayRequest = new EquitySecurityGatewayRequest(equitySecurityNewRequest);
}
securityGatewayRequest.setListing(listingRequests);
securityGatewayRequest.setCouponPeriods(couponPeriodRequests);
securityGatewayRequest.setFixedIncomesCashFlow(cashFlowRequests);
kafkaSender.sendRequestToQueue(Consts.DESTINATION_SECURITY_MULTIREQUEST, securityGatewayRequest);
}
}
private <T extends WithSecurityId> List<T> groupByCompanyId(List<T> listToProcess, UUID securityId) {
return listToProcess.stream().filter(t -> t.getSecurityId().equals(securityId)).collect(Collectors.toList());
issueCompanyProcessor.sendRequest(request);
securityService.sendRequest(request);
}
}

View file

@ -0,0 +1,69 @@
package ru.spcex.clearing.gatewayapi.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.gatewayapi.controller.request.WithSecurityId;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.*;
import ru.spcex.clearing.gatewayapi.service.adapter.SecurityFondRequestAdapter;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.company.EquitySecurityGatewayRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.FixedIncomeGatewayRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.SecurityFondGatewayRequest;
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 java.util.List;
import java.util.UUID;
import java.util.stream.Collectors;
@Service
public class SecurityService {
private final Logger log = LoggerFactory.getLogger(getClass());
private final SecurityFondRequestAdapter securityRequestAdapter;
private final KafkaSender kafkaSender;
public SecurityService(SecurityFondRequestAdapter securityRequestAdapter,
KafkaSender kafkaSender) {
this.securityRequestAdapter = securityRequestAdapter;
this.kafkaSender = kafkaSender;
}
public void sendRequest(FondListingsRequest request) {
List<FondSecurity> securityList = request.getSecurities();
for (FondSecurity security : securityList) {
UUID securityId = security.getId();
log.debug("Grouping by securityId: {}", securityId);
List<IncomeListing> incomeListings = groupByCompanyId(request.getListingList(), securityId);
List<CouponSchedule> couponSchedules = groupByCompanyId(request.getCouponSchedules(), securityId);
List<Nominal> nominals = groupByCompanyId(request.getNominalList(), securityId);
List<ListingNewRequest> listingRequests = incomeListings.stream().
map(securityRequestAdapter::toIncomeListing).toList();
List<CouponPeriodNewRequest> couponPeriodRequests = couponSchedules.stream().
map(securityRequestAdapter::toCouponPeriodRequest).toList();
List<FixedIncomeCashFlowNewRequest> cashFlowRequests = nominals.stream().
map(securityRequestAdapter::toFixedIncomeCashFlowRequest).toList();
SecurityFondGatewayRequest<?> securityGatewayRequest;
if (InstrumentType.BOND.equalsByKey(security.getInstrumentType())) {
FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest = securityRequestAdapter.toFixedIncomeSecurityRequest(security);
securityGatewayRequest = new FixedIncomeGatewayRequest(fixedIncomeSecurityNewRequest);
} else {
EquitySecurityNewRequest equitySecurityNewRequest = securityRequestAdapter.toEquitySecurityNewRequest(security);
securityGatewayRequest = new EquitySecurityGatewayRequest(equitySecurityNewRequest);
}
securityGatewayRequest.setListing(listingRequests);
securityGatewayRequest.setCouponPeriods(couponPeriodRequests);
securityGatewayRequest.setFixedIncomesCashFlow(cashFlowRequests);
kafkaSender.sendRequestToQueue(Consts.DESTINATION_SECURITY_MULTIREQUEST, securityGatewayRequest);
}
}
private <T extends WithSecurityId> List<T> groupByCompanyId(List<T> listToProcess, UUID securityId) {
return listToProcess.stream().filter(t -> t.getSecurityId().equals(securityId)).collect(Collectors.toList());
}
}

View file

@ -8,7 +8,7 @@ import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.Nominal;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.*;
@Service
public class SecurityRequestAdapter {
public class SecurityFondRequestAdapter {
public FixedIncomeSecurityNewRequest toFixedIncomeSecurityRequest(FondSecurity fondSecurity) {
FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest = new FixedIncomeSecurityNewRequest();

View file

@ -0,0 +1,52 @@
package ru.spcex.clearing.gatewayapi.service.adapter;
import org.springframework.stereotype.Service;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompany;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompanyInfo;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompanySymbols;
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerContact;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyInfoUpdateRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanySymbolNewRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest;
import ru.spcex.platform.enumeration.CompanySymbol;
@Service
public class SecurityMkrRequestAdapter {
public CompanyNewRequest toCompanyRequest(IssuerCompany issuerCompany) {
CompanyNewRequest companyNewRequest = new CompanyNewRequest();
companyNewRequest.setShortName(issuerCompany.getShortName());
companyNewRequest.setFullName(issuerCompany.getFullName());
companyNewRequest.setWorkflowStatus(issuerCompany.getWorkflowStatus());
companyNewRequest.setCompanySymbol(CompanySymbol.UUID.getKey());
companyNewRequest.setCompanySymbolValue(issuerCompany.getId().toString());
return companyNewRequest;
}
public CompanyInfoUpdateRequest toCompanyInfoRequest(IssuerCompanyInfo issuerCompanyInfo) {
CompanyInfoUpdateRequest companyInfoUpdateRequest = new CompanyInfoUpdateRequest();
companyInfoUpdateRequest.setCountryCode(issuerCompanyInfo.getCountryCode());
companyInfoUpdateRequest.setLegalKind(issuerCompanyInfo.getLegalKind());
companyInfoUpdateRequest.setOrganizationType(issuerCompanyInfo.getOrganizationType());
companyInfoUpdateRequest.setResidence(issuerCompanyInfo.getResidence());
companyInfoUpdateRequest.setShortNameEng(issuerCompanyInfo.getShortNameEng());
companyInfoUpdateRequest.setFullNameEng(issuerCompanyInfo.getFullNameEng());
return companyInfoUpdateRequest;
}
public CompanySymbolNewRequest toCompanySymbolRequest(IssuerCompanySymbols issuerCompanySymbols) {
CompanySymbolNewRequest companySymbolNewRequest = new CompanySymbolNewRequest();
companySymbolNewRequest.setCompanySymbol(issuerCompanySymbols.getCompanySymbol());
companySymbolNewRequest.setCompanySymbolValue(issuerCompanySymbols.getCompanySymbolValue());
return companySymbolNewRequest;
}
public ContactNewRequest toContactRequest(IssuerContact issuerContact) {
ContactNewRequest contactNewRequest = new ContactNewRequest();
contactNewRequest.setContactType(issuerContact.getContactType());
contactNewRequest.setContactValue(issuerContact.getContactValue());
return contactNewRequest;
}
}

View file

@ -1,176 +0,0 @@
package ru.spcex.clearing.gatewayapi.logic.listings_fond;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mockito;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.spcex.clearing.gatewayapi.config.ProcessorConfiguration;
import ru.spcex.clearing.gatewayapi.config.TestConfiguration;
import ru.spcex.clearing.gatewayapi.logic.Processor;
import ru.spcex.clearing.gatewayapi.logic.Stage;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.company.CompanyNewRequest;
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.service.sender.KafkaSender;
import ru.spcex.clearing.test.config.ImdgTestConfig;
import ru.spcex.clearing.test.config.KafkaTestConfig;
import ru.spcex.platform.enumeration.CompanySymbol;
import ru.spcex.platform.enumeration.InstrumentType;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
// Суть теста: прогоняем через pipeline запрос и в mockSender ловим переданные company и security сервисам запросы
@ExtendWith(SpringExtension.class)
@ContextConfiguration(classes = {
ProcessorConfiguration.class,
TestConfiguration.class,
ImdgTestConfig.class,
KafkaTestConfig.class})
class ListingFondRequestTest implements InitializingBean {
// В этот объект будут помещаться отправленные сообщения
final Map<String, List<Object>> sentMessages = new HashMap<>();
@Autowired
@Qualifier("testObjectMapper")
ObjectMapper objectMapper;
@Autowired
Processor<FondListingsRequestParam> processor;
@Autowired
ImdgProvider imdgProvider;
KafkaSender mockSender = Mockito.mock(KafkaSender.class);
Imdg<Company> companyImdg;
Imdg<CompanySymbols> companySymbolsImdg;
Imdg<FixedIncomeSecurity> fixedIncomeSecurityImdg;
Imdg<EquitySecurity> equitySecurityImdg;
// @Test
// public void test() throws IOException {
// File requestFile = new File(getClass().getClassLoader().getResource("requests/testFondSecuritiesRequest_1.json").getFile());
//
// FondListingsRequest fondListingsRequest = objectMapper.readValue(requestFile, FondListingsRequest.class);
// FondListingsRequestParam param = new FondListingsRequestParam(fondListingsRequest);
//
// ProcessResult result = processor.process(param);
//
// // Проверяем только количество запросов, поскольку выполнение этих запросов зависит от security и company модулей
// // todo upd: поскольку будет производиться переход к полностью асинхронному запросам, этот тест возможно стоит тоже поправить
// Assertions.assertNull(result.getError());
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_COMPANY_NEW).size(), 1);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_COMPANY_INFO_UPDATE).size(), 2);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_COMPANY_UPDATE).size(), 1);
// Assertions.assertEquals(sentMessages.get(Consts.LISTING_NEW).size(), 4);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_COUPON_PERIOD_NEW).size(), 1);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_COMPANY_SYMBOL_NEW).size(), 2);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_COMPANY_SYMBOL_UPDATE).size(), 2);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_CONTACT_NEW).size(), 2);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_FIXED_INCOME_SECURITY_UPDATE).size(), 2);
// Assertions.assertEquals(sentMessages.get(Consts.DESTINATION_EQUITY_SECURITY_UPDATE).size(), 2);
// }
@Override
public void afterPropertiesSet() throws Exception {
companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
companySymbolsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class);
Company existCompany1 = new Company();
CompanySymbols existCompany1Symbols = new CompanySymbols();
Long company1Id = companyImdg.insert(existCompany1);
existCompany1Symbols.setCompanySymbol(CompanySymbol.UUID.getKey());
existCompany1Symbols.setCompanySymbolValue("ef47f76c-bcea-11ed-afa1-0242ac120000");
existCompany1Symbols.setCompanyId(company1Id);
companySymbolsImdg.insert(existCompany1Symbols);
Company existCompany2 = new Company();
CompanySymbols existCompany2Symbols = new CompanySymbols();
Long company2Id = companyImdg.insert(existCompany2);
existCompany2Symbols.setCompanySymbol(CompanySymbol.UUID.getKey());
existCompany2Symbols.setCompanySymbolValue("ef47f76c-bcea-11ed-afa1-0242ac120001");
existCompany2Symbols.setCompanyId(company2Id);
companySymbolsImdg.insert(existCompany2Symbols);
FixedIncomeSecurity fixedIncomeSecurity_1 = new FixedIncomeSecurity();
fixedIncomeSecurity_1.setInstrumentType(InstrumentType.BOND.getKey());
fixedIncomeSecurity_1.setSecuritySymbol("VTB_001P-06 TEST 3");
fixedIncomeSecurity_1.setWorkflowStatus(WorkflowStatus.Active.getKey());
fixedIncomeSecurity_1.setShortName("Биржевые облигации ВТБ TEST 3");
fixedIncomeSecurityImdg.insert(fixedIncomeSecurity_1);
EquitySecurity equitySecurity_1 = new EquitySecurity();
equitySecurity_1.setInstrumentType(InstrumentType.EQTY.getKey());
equitySecurity_1.setSecuritySymbol("VTB_001P-06 TEST 4");
equitySecurity_1.setWorkflowStatus(WorkflowStatus.Active.getKey());
equitySecurity_1.setShortName("Биржевые облигации ВТБ TEST 4");
equitySecurityImdg.insert(equitySecurity_1);
// мок kafkaSender, чтобы он лишь сохранял в мапу переданные аргументы
Mockito
.when(mockSender.sendRequestToQueue(Mockito.any(), Mockito.any()))
.thenAnswer(
invocation -> {
String destination = invocation.getArgument(0);
Object invocationArgument = invocation.getArgument(1);
if (invocationArgument instanceof CompanyNewRequest newRequest) {
Company company = new Company();
company.setShortName(newRequest.getShortName());
company.setFullName(newRequest.getFullName());
company.setWorkflowStatus(newRequest.getWorkflowStatus());
companyImdg.insert(company);
}
if (invocationArgument instanceof FixedIncomeSecurityNewRequest newRequest) {
FixedIncomeSecurity fixedIncomeSecurity = new FixedIncomeSecurity();
fixedIncomeSecurity.setInstrumentType(newRequest.getInstrumentType());
fixedIncomeSecurity.setSecuritySymbol(newRequest.getSecuritySymbol());
fixedIncomeSecurity.setWorkflowStatus(newRequest.getWorkflowStatus());
fixedIncomeSecurity.setShortName(newRequest.getShortName());
fixedIncomeSecurityImdg.insert(fixedIncomeSecurity);
}
if (invocationArgument instanceof EquitySecurityNewRequest newRequest) {
EquitySecurity equitySecurity = new EquitySecurity();
equitySecurity.setInstrumentType(newRequest.getInstrumentType());
equitySecurity.setSecuritySymbol(newRequest.getSecuritySymbol());
equitySecurity.setWorkflowStatus(newRequest.getWorkflowStatus());
equitySecurity.setShortName(newRequest.getShortName());
equitySecurityImdg.insert(equitySecurity);
}
List<Object> currentMessageList = sentMessages.computeIfAbsent(destination, k -> new ArrayList<>());
currentMessageList.add(invocationArgument);
return 0L;
}
);
// в pipeline для processor заменяем стадии с отправкой сообщений
List<Stage<FondListingsRequestParam>> pipeline = processor.getPipeline();
for (Stage<FondListingsRequestParam> stage : pipeline) {
if (stage instanceof SendMessageToSecurityServiceWithFondSecurities) {
int index = pipeline.indexOf(stage);
Stage<FondListingsRequestParam> mockedStage = new SendMessageToSecurityServiceWithFondSecurities(mockSender, imdgProvider);
pipeline.set(index, mockedStage);
}
if (stage instanceof SendMessageToCompanyServiceWithIssuerCompanies) {
int index = pipeline.indexOf(stage);
Stage<FondListingsRequestParam> mockedStage = new SendMessageToCompanyServiceWithIssuerCompanies(mockSender, imdgProvider);
pipeline.set(index, mockedStage);
}
}
}
}

View file

@ -43,7 +43,8 @@ public interface Consts {
String DESTINATION_PLANNER_UPDATE = "planner-update";
String DESTINATION_PLANNER_DELETE = "planner-delete";
String DESTINATION_COMPANY_MULTIREQUEST = "company-multirequest-new";
String DESTINATION_COMPANY_GATEWAY_REQUEST = "company-gateway-request";
String DESTINATION_ISSUER_COMPANY_GATEWAY_REQUEST = "company-issuer-company-gateway-request";
String DESTINATION_SECURITY_MULTIREQUEST = "company-multirequest-new";
String DESTINATION_COMPANY_NEW = "company-new";
String DESTINATION_COMPANY_DELETE = "company-delete";

View file

@ -2,7 +2,7 @@ package ru.spcex.clearing.platform.messaging.domain.cud.company;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.EquitySecurityNewRequest;
public class EquitySecurityGatewayRequest extends SecurityGatewayRequest<EquitySecurityNewRequest> {
public class EquitySecurityGatewayRequest extends SecurityFondGatewayRequest<EquitySecurityNewRequest> {
public EquitySecurityGatewayRequest(EquitySecurityNewRequest request) {
super(request);
}

View file

@ -2,7 +2,7 @@ package ru.spcex.clearing.platform.messaging.domain.cud.company;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.FixedIncomeSecurityNewRequest;
public class FixedIncomeGatewayRequest extends SecurityGatewayRequest<FixedIncomeSecurityNewRequest> {
public class FixedIncomeGatewayRequest extends SecurityFondGatewayRequest<FixedIncomeSecurityNewRequest> {
public FixedIncomeGatewayRequest(FixedIncomeSecurityNewRequest request) {
super(request);
}

View file

@ -10,7 +10,7 @@ import ru.spcex.platform.classes.base.interfaces.WithSecuritySymbol;
import java.util.ArrayList;
import java.util.List;
public abstract class SecurityGatewayRequest<T extends WithSecuritySymbol & WithInstrumentType> {
public abstract class SecurityFondGatewayRequest<T extends WithSecuritySymbol & WithInstrumentType> {
@JsonProperty
private String uuid;
@ -23,7 +23,7 @@ public abstract class SecurityGatewayRequest<T extends WithSecuritySymbol & With
@JsonProperty
private List<FixedIncomeCashFlowNewRequest> fixedIncomesCashFlow = new ArrayList<>();
public SecurityGatewayRequest(T security) {
public SecurityFondGatewayRequest(T security) {
this.security = security;
}

View file

@ -0,0 +1,58 @@
package ru.spcex.clearing.platform.messaging.domain.cud.company;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.util.ArrayList;
import java.util.List;
public class SecurityMkrGatewayRequest {
@JsonProperty
private String uuid;
@JsonProperty
private CompanyNewRequest companyNewRequest;
private CompanyInfoUpdateRequest companyInfoUpdateRequest;
@JsonProperty
private List<CompanySymbolNewRequest> companySymbolNewRequests = new ArrayList<>();
@JsonProperty
private List<ContactNewRequest> contactNewRequests = new ArrayList<>();
public String getUuid() {
return uuid;
}
public void setUuid(String uuid) {
this.uuid = uuid;
}
public CompanyNewRequest getCompanyNewRequest() {
return companyNewRequest;
}
public void setCompanyNewRequest(CompanyNewRequest companyNewRequest) {
this.companyNewRequest = companyNewRequest;
}
public CompanyInfoUpdateRequest getCompanyInfoUpdateRequest() {
return companyInfoUpdateRequest;
}
public void setCompanyInfoUpdateRequest(CompanyInfoUpdateRequest companyInfoUpdateRequest) {
this.companyInfoUpdateRequest = companyInfoUpdateRequest;
}
public List<CompanySymbolNewRequest> getCompanySymbolNewRequests() {
return companySymbolNewRequests;
}
public void setCompanySymbolNewRequests(List<CompanySymbolNewRequest> companySymbolNewRequests) {
this.companySymbolNewRequests = companySymbolNewRequests;
}
public List<ContactNewRequest> getContactNewRequests() {
return contactNewRequests;
}
public void setContactNewRequests(List<ContactNewRequest> contactNewRequests) {
this.contactNewRequests = contactNewRequests;
}
}