some refactroing gateway
This commit is contained in:
parent
c15f3fba9d
commit
b6f94a4069
17 changed files with 99 additions and 334 deletions
|
|
@ -1,32 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.config;
|
|
||||||
|
|
||||||
import org.springframework.beans.factory.annotation.Qualifier;
|
|
||||||
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_mm.MMListingsRequestParam;
|
|
||||||
import ru.spcex.clearing.gatewayapi.logic.listings_mm.PrepareExchangeInstruments;
|
|
||||||
import ru.spcex.clearing.gatewayapi.logic.listings_mm.SendMessageToSecurityServiceWithMMSecurities;
|
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
|
||||||
|
|
||||||
import java.util.ArrayList;
|
|
||||||
import java.util.List;
|
|
||||||
|
|
||||||
@Configuration
|
|
||||||
public class ProcessorConfiguration {
|
|
||||||
|
|
||||||
|
|
||||||
@Qualifier("mmListingsRequestProcessor")
|
|
||||||
@Bean
|
|
||||||
public Processor<MMListingsRequestParam> mmListingsRequestProcessor(ImdgProvider imdgProvider, KafkaSender kafkaSender) {
|
|
||||||
Processor<MMListingsRequestParam> processor = new Processor<>();
|
|
||||||
List<Stage<MMListingsRequestParam>> pipeline = new ArrayList<>();
|
|
||||||
pipeline.add(new PrepareExchangeInstruments());
|
|
||||||
pipeline.add(new SendMessageToSecurityServiceWithMMSecurities(kafkaSender));
|
|
||||||
processor.setPipeline(pipeline);
|
|
||||||
return processor;
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
@ -17,11 +17,9 @@ import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondListings
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.MMListingsRequest;
|
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.MMListingsRequest;
|
||||||
import ru.spcex.clearing.gatewayapi.controller.response.CommonResponse;
|
import ru.spcex.clearing.gatewayapi.controller.response.CommonResponse;
|
||||||
import ru.spcex.clearing.gatewayapi.exception.GatewayException;
|
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_mm.MMListingsRequestParam;
|
|
||||||
import ru.spcex.clearing.gatewayapi.service.CompanyProcessor;
|
import ru.spcex.clearing.gatewayapi.service.CompanyProcessor;
|
||||||
import ru.spcex.clearing.gatewayapi.service.ListingFondProcessor;
|
import ru.spcex.clearing.gatewayapi.service.ListingFondProcessor;
|
||||||
|
import ru.spcex.clearing.gatewayapi.service.ListingMMProcessor;
|
||||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
||||||
|
|
@ -37,23 +35,22 @@ public class GatewayController {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
private final IMessageResolver messageResolver;
|
private final IMessageResolver messageResolver;
|
||||||
private final Processor<MMListingsRequestParam> mmListingsRequestProcessor;
|
|
||||||
private final CompanyProcessor companiesProcessor;
|
private final CompanyProcessor companiesProcessor;
|
||||||
private final ListingFondProcessor listingFondProcessor;
|
private final ListingFondProcessor listingFondProcessor;
|
||||||
|
private final ListingMMProcessor listingMMProcessor;
|
||||||
private final ExecutorService executor;
|
private final ExecutorService executor;
|
||||||
|
|
||||||
private List<String> validTypes = Arrays.asList("DAY_START", "ON_DEMAND");
|
private List<String> validTypes = Arrays.asList("DAY_START", "ON_DEMAND");
|
||||||
|
|
||||||
public GatewayController(
|
public GatewayController(IMessageResolver messageResolver,
|
||||||
IMessageResolver messageResolver,
|
CompanyProcessor companiesProcessor,
|
||||||
@Qualifier("mmListingsRequestProcessor") Processor<MMListingsRequestParam> mmListingsRequestProcessor,
|
ListingFondProcessor listingFondProcessor,
|
||||||
CompanyProcessor companiesProcessor,
|
ListingMMProcessor listingMMProcessor,
|
||||||
ListingFondProcessor listingFondProcessor,
|
@Qualifier("gatewayExecutor") ExecutorService executor) {
|
||||||
@Qualifier("gatewayExecutor") ExecutorService executor) {
|
|
||||||
this.messageResolver = messageResolver;
|
this.messageResolver = messageResolver;
|
||||||
this.mmListingsRequestProcessor = mmListingsRequestProcessor;
|
|
||||||
this.companiesProcessor = companiesProcessor;
|
this.companiesProcessor = companiesProcessor;
|
||||||
this.listingFondProcessor = listingFondProcessor;
|
this.listingFondProcessor = listingFondProcessor;
|
||||||
|
this.listingMMProcessor = listingMMProcessor;
|
||||||
this.executor = executor;
|
this.executor = executor;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -91,15 +88,7 @@ public class GatewayController {
|
||||||
@ResponseBody
|
@ResponseBody
|
||||||
public CommonResponse listingsMM(@RequestBody MMListingsRequest request) {
|
public CommonResponse listingsMM(@RequestBody MMListingsRequest request) {
|
||||||
// request.validate(validTypes, List.of(Section.MKR));
|
// request.validate(validTypes, List.of(Section.MKR));
|
||||||
MMListingsRequestParam requestParam = new MMListingsRequestParam(request);
|
executor.submit(() -> listingMMProcessor.process(request));
|
||||||
executor.submit(() -> {
|
|
||||||
ProcessResult processResult = mmListingsRequestProcessor.process(requestParam);
|
|
||||||
if (processResult.getError() != null) {
|
|
||||||
log.warn("Exception while process listing_mm: {}", processResult.getError());
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
|
|
||||||
return createResponse(request, true);
|
return createResponse(request, true);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,24 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.logic;
|
|
||||||
|
|
||||||
import ru.spcex.clearing.gatewayapi.errors.GatewayError;
|
|
||||||
|
|
||||||
public class ProcessResult {
|
|
||||||
private GatewayError error = null;
|
|
||||||
private Object[] errorArgs;
|
|
||||||
|
|
||||||
public ProcessResult() {
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
public ProcessResult(GatewayError error, Object ... errorArgs) {
|
|
||||||
this.error = error;
|
|
||||||
}
|
|
||||||
|
|
||||||
public GatewayError getError() {
|
|
||||||
return error;
|
|
||||||
}
|
|
||||||
|
|
||||||
public void setError(GatewayError error) {
|
|
||||||
this.error = error;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -1,34 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.logic;
|
|
||||||
|
|
||||||
import org.apache.commons.lang3.exception.ExceptionUtils;
|
|
||||||
import org.slf4j.Logger;
|
|
||||||
import org.slf4j.LoggerFactory;
|
|
||||||
import ru.spcex.clearing.gatewayapi.errors.GatewayError;
|
|
||||||
|
|
||||||
import java.util.List;
|
|
||||||
|
|
||||||
public class Processor<T> {
|
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
|
||||||
private List<Stage<T>> pipeline;
|
|
||||||
|
|
||||||
public ProcessResult process(T param) {
|
|
||||||
for (Stage<T> stage : pipeline) {
|
|
||||||
try {
|
|
||||||
ProcessResult processResult = stage.process(param);
|
|
||||||
if (processResult != null) return processResult;
|
|
||||||
} catch (RuntimeException e) {
|
|
||||||
log.warn("Stage : {} have not been complete successful: {}", stage.getClass(), ExceptionUtils.getMessage(e));
|
|
||||||
return new ProcessResult(GatewayError.InternalError);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return new ProcessResult();
|
|
||||||
}
|
|
||||||
|
|
||||||
public List<Stage<T>> getPipeline() {
|
|
||||||
return pipeline;
|
|
||||||
}
|
|
||||||
|
|
||||||
public void setPipeline(List<Stage<T>> pipeline) {
|
|
||||||
this.pipeline = pipeline;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -1,5 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.logic;
|
|
||||||
|
|
||||||
public abstract class Stage<T> {
|
|
||||||
public abstract ProcessResult process(T param);
|
|
||||||
}
|
|
||||||
|
|
@ -1,28 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.logic.listings_mm;
|
|
||||||
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.ExchangeInstrument;
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.MMListingsRequest;
|
|
||||||
|
|
||||||
import java.util.List;
|
|
||||||
|
|
||||||
public class MMListingsRequestParam {
|
|
||||||
private final MMListingsRequest request;
|
|
||||||
private List<ExchangeInstrument> exchangeInstrumentsList;
|
|
||||||
|
|
||||||
public MMListingsRequestParam(MMListingsRequest request) {
|
|
||||||
this.request = request;
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
public MMListingsRequest getRequest() {
|
|
||||||
return request;
|
|
||||||
}
|
|
||||||
|
|
||||||
public List<ExchangeInstrument> getExchangeInstrumentsList() {
|
|
||||||
return exchangeInstrumentsList;
|
|
||||||
}
|
|
||||||
|
|
||||||
public void setExchangeInstrumentsList(List<ExchangeInstrument> exchangeInstrumentsList) {
|
|
||||||
this.exchangeInstrumentsList = exchangeInstrumentsList;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -1,108 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.logic.listings_mm;
|
|
||||||
|
|
||||||
import org.slf4j.Logger;
|
|
||||||
import org.slf4j.LoggerFactory;
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.DirAuctionBiddingType;
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.DirTradingModeMkr;
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.ExchangeInstrument;
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.TradingMode;
|
|
||||||
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 PrepareExchangeInstruments extends Stage<MMListingsRequestParam> {
|
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public ProcessResult process(MMListingsRequestParam param) {
|
|
||||||
List<ExchangeInstrument> exchangeInstrumentList = param.getExchangeInstrumentsList();
|
|
||||||
|
|
||||||
// Проверяем что в exchange_instrument нет объектов с одинаковым id
|
|
||||||
Map<UUID, ExchangeInstrument> exchangeInstrumentMap = new HashMap<>();
|
|
||||||
for (ExchangeInstrument exchangeInstrument : exchangeInstrumentList) {
|
|
||||||
if (exchangeInstrumentMap.containsKey(exchangeInstrument.getId())) {
|
|
||||||
log.warn("For exchange_instrument.UUID {} found > 1 element from request, use first", exchangeInstrument.getId());
|
|
||||||
exchangeInstrument.setInvalidData(true);
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
exchangeInstrumentMap.put(exchangeInstrument.getId(), exchangeInstrument);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Находим для exchange_instrument соответствующий trading_mode по условию:
|
|
||||||
// exchange_instrument.id = trading_mode.exchange_instrument_id
|
|
||||||
List<TradingMode> tradingModeList = param.getRequest().getTradingModeList();
|
|
||||||
Map<UUID, TradingMode> tradingModeMap = new HashMap<>();
|
|
||||||
for (TradingMode tradingMode : tradingModeList) {
|
|
||||||
if (tradingModeMap.containsKey(tradingMode.getTradingModeId())) {
|
|
||||||
log.warn("For trading_mode.UUID {} found > 1 element from request, use first", tradingMode.getTradingModeId());
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
tradingModeMap.put(tradingMode.getTradingModeId(), tradingMode);
|
|
||||||
UUID exchangeInstrumentId = tradingMode.getExchangeInstrumentId();
|
|
||||||
ExchangeInstrument exchangeInstrument = exchangeInstrumentMap.get(exchangeInstrumentId);
|
|
||||||
if (exchangeInstrument != null) {
|
|
||||||
if (exchangeInstrument.getTradingMode() != null) {
|
|
||||||
log.warn("For exchange_instrument.id {} found > 1 trading_mode, use first", tradingMode.getExchangeInstrumentId());
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
exchangeInstrument.setTradingMode(tradingMode);
|
|
||||||
} else {
|
|
||||||
log.warn("For trading_mode.exchange_instrument_id {} not found exchange_instrument, trading_mode skipped", exchangeInstrumentId);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Находим для exchange_instrument соответсвующий dir_auction_bidding_type по условию:
|
|
||||||
// exchange_instrument.auction_bidding_type_id = dir_auction_bidding_type.id
|
|
||||||
Map<UUID, DirAuctionBiddingType> dirAuctionBiddingTypeMap = new HashMap<>();
|
|
||||||
for (DirAuctionBiddingType dirAuctionBiddingType : param.getRequest().getDirAuctionBiddingTypeList()) {
|
|
||||||
if (dirAuctionBiddingTypeMap.containsKey(dirAuctionBiddingType.getId())) {
|
|
||||||
log.warn("For dir_auction_bidding_type.id {} found > 1 element from request, use first", dirAuctionBiddingType.getId());
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
dirAuctionBiddingTypeMap.put(dirAuctionBiddingType.getId(), dirAuctionBiddingType);
|
|
||||||
}
|
|
||||||
for (ExchangeInstrument exchangeInstrument : exchangeInstrumentList) {
|
|
||||||
if (exchangeInstrument.isInvalidData()) continue;
|
|
||||||
DirAuctionBiddingType dirAuctionBiddingType = dirAuctionBiddingTypeMap.get(exchangeInstrument.getAuctionBiddingTypeId());
|
|
||||||
if (dirAuctionBiddingType == null) {
|
|
||||||
log.warn("For exchange_instrument.id {} not found dir_auction_bidding_type", exchangeInstrument.getId());
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
exchangeInstrument.setDirAuctionBiddingType(dirAuctionBiddingType);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Находим для trading_mode соответствующий dir_trading_mode_mkr по условию:
|
|
||||||
// trading_mode.trading_mode_id = dir_trading_mode_mkr.id
|
|
||||||
List<DirTradingModeMkr> dirTradingModeMkrList = param.getRequest().getDirTradingModeMkrList();
|
|
||||||
for (DirTradingModeMkr dirTradingModeMkr : dirTradingModeMkrList) {
|
|
||||||
UUID tradingModeUUID = dirTradingModeMkr.getId();
|
|
||||||
TradingMode tradingMode = tradingModeMap.get(tradingModeUUID);
|
|
||||||
if (tradingMode != null) {
|
|
||||||
if (tradingMode.getDirTradingModeMkr() != null) {
|
|
||||||
log.warn("For trading_mode.trading_more_id {} found > 1 dir_trading_mode_mkr, use first", tradingModeUUID);
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
tradingMode.setDirTradingModeMkr(dirTradingModeMkr);
|
|
||||||
} else {
|
|
||||||
log.warn("For dir_trading_mode_mkr.id {} not found trading_mode (for trading_mode_id), dir_trading_mode skipped", tradingModeUUID);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
List<ExchangeInstrument> filteredExchangeInstrument = exchangeInstrumentList.stream()
|
|
||||||
.filter(exchangeInstrument -> !exchangeInstrument.isInvalidData())
|
|
||||||
.collect(Collectors.toList());
|
|
||||||
|
|
||||||
if (filteredExchangeInstrument.isEmpty()) {
|
|
||||||
return new ProcessResult(GatewayError.SecurityNotFound, exchangeInstrumentList.stream().map(fondSecurity -> fondSecurity.getId().toString()).collect(Collectors.joining(", ")));
|
|
||||||
}
|
|
||||||
|
|
||||||
param.setExchangeInstrumentsList(filteredExchangeInstrument);
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -1,58 +0,0 @@
|
||||||
package ru.spcex.clearing.gatewayapi.logic.listings_mm;
|
|
||||||
|
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.ExchangeInstrument;
|
|
||||||
import ru.spcex.clearing.gatewayapi.logic.ProcessResult;
|
|
||||||
import ru.spcex.clearing.gatewayapi.logic.Stage;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest;
|
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
|
||||||
import ru.spcex.platform.enumeration.InstrumentType;
|
|
||||||
|
|
||||||
import java.util.List;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Отправка сообщения к company-service на создание/обновление company
|
|
||||||
* todo Пока убрал совсем не рабочую логику с sendToQueueWaitForAnswer
|
|
||||||
*/
|
|
||||||
public class SendMessageToSecurityServiceWithMMSecurities extends Stage<MMListingsRequestParam> {
|
|
||||||
private final KafkaSender kafkaSender;
|
|
||||||
|
|
||||||
public SendMessageToSecurityServiceWithMMSecurities(KafkaSender kafkaSender) {
|
|
||||||
this.kafkaSender = kafkaSender;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public ProcessResult process(MMListingsRequestParam param) {
|
|
||||||
List<ExchangeInstrument> exchangeInstrumentList = param.getExchangeInstrumentsList();
|
|
||||||
for (ExchangeInstrument exchangeInstrument : exchangeInstrumentList) {
|
|
||||||
if (exchangeInstrument.isInvalidData()) continue;
|
|
||||||
|
|
||||||
if (!exchangeInstrument.isAlreadyExist()) {
|
|
||||||
MoneyMarketSecurityNewRequest moneyMarketSecurityNewRequest = new MoneyMarketSecurityNewRequest();
|
|
||||||
moneyMarketSecurityNewRequest.setNominalCurrency(exchangeInstrument.getSettlementCurrencyLetterCode());
|
|
||||||
moneyMarketSecurityNewRequest.setInstrumentType(InstrumentType.RATE.getKey());
|
|
||||||
moneyMarketSecurityNewRequest.setShortName(exchangeInstrument.getName());
|
|
||||||
moneyMarketSecurityNewRequest.setSecuritySymbol(exchangeInstrument.getSecuritySymbol());
|
|
||||||
moneyMarketSecurityNewRequest.setLotSize(exchangeInstrument.getLot());
|
|
||||||
moneyMarketSecurityNewRequest.setTermType(exchangeInstrument.getBankDepositAgreementType());
|
|
||||||
moneyMarketSecurityNewRequest.setIssuerId(exchangeInstrument.getIssuerId());
|
|
||||||
moneyMarketSecurityNewRequest.setWorkflowStatus(exchangeInstrument.getWorkflowStatus());
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, moneyMarketSecurityNewRequest);
|
|
||||||
} else {
|
|
||||||
MoneyMarketSecurityUpdateRequest moneyMarketSecurityUpdateRequest = new MoneyMarketSecurityUpdateRequest();
|
|
||||||
moneyMarketSecurityUpdateRequest.setId(exchangeInstrument.getMapId());
|
|
||||||
moneyMarketSecurityUpdateRequest.setNominalCurrency(exchangeInstrument.getSettlementCurrencyLetterCode());
|
|
||||||
moneyMarketSecurityUpdateRequest.setInstrumentType(InstrumentType.RATE.getKey());
|
|
||||||
moneyMarketSecurityUpdateRequest.setShortName(exchangeInstrument.getName());
|
|
||||||
moneyMarketSecurityUpdateRequest.setSecuritySymbol(exchangeInstrument.getSecuritySymbol());
|
|
||||||
moneyMarketSecurityUpdateRequest.setLotSize(exchangeInstrument.getLot());
|
|
||||||
moneyMarketSecurityUpdateRequest.setTermType(exchangeInstrument.getBankDepositAgreementType());
|
|
||||||
moneyMarketSecurityUpdateRequest.setIssuerId(exchangeInstrument.getIssuerId());
|
|
||||||
moneyMarketSecurityUpdateRequest.setWorkflowStatus(exchangeInstrument.getWorkflowStatus());
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW, moneyMarketSecurityUpdateRequest);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -147,8 +147,8 @@ public class GatewayService extends QueueConsumer implements InitializingBean {
|
||||||
"amount", assetOperationRequest.getAmount(),
|
"amount", assetOperationRequest.getAmount(),
|
||||||
"quantity", assetOperationRequest.getQuantity(),
|
"quantity", assetOperationRequest.getQuantity(),
|
||||||
"code", assetOperationRequest.getCode(),
|
"code", assetOperationRequest.getCode(),
|
||||||
"asset", assetOperationRequest.getAsset(),
|
"asset", assetOperationRequest.getSecuritySymbol(),
|
||||||
"firm_id", assetOperationRequest.getFirmId());
|
"firm_id", assetOperationRequest.getTradingCode());
|
||||||
|
|
||||||
OutboundRequest outboundRequest = OutboundRequestBuilder.builder()
|
OutboundRequest outboundRequest = OutboundRequestBuilder.builder()
|
||||||
.section(Section.MKR.getKey())
|
.section(Section.MKR.getKey())
|
||||||
|
|
|
||||||
|
|
@ -8,7 +8,7 @@ import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.FondListings
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerCompany;
|
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.IssuerCompanySymbols;
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerContact;
|
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.IssuerContact;
|
||||||
import ru.spcex.clearing.gatewayapi.service.adapter.SecurityMkrRequestAdapter;
|
import ru.spcex.clearing.gatewayapi.service.adapter.IssuerCompanyRequestAdapter;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
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.CompanySymbolNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest;
|
||||||
|
|
@ -23,12 +23,12 @@ import java.util.stream.Collectors;
|
||||||
public class IssueCompanyService {
|
public class IssueCompanyService {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
private final SecurityMkrRequestAdapter securityMkrRequestAdapter;
|
private final IssuerCompanyRequestAdapter issuerCompanyRequestAdapter;
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
|
|
||||||
public IssueCompanyService(SecurityMkrRequestAdapter securityMkrRequestAdapter,
|
public IssueCompanyService(IssuerCompanyRequestAdapter issuerCompanyRequestAdapter,
|
||||||
KafkaSender kafkaSender) {
|
KafkaSender kafkaSender) {
|
||||||
this.securityMkrRequestAdapter = securityMkrRequestAdapter;
|
this.issuerCompanyRequestAdapter = issuerCompanyRequestAdapter;
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -41,18 +41,18 @@ public class IssueCompanyService {
|
||||||
List<IssuerContact> issuerContactList = groupByCompanyId(request.getIssuerContactList(), companyId);
|
List<IssuerContact> issuerContactList = groupByCompanyId(request.getIssuerContactList(), companyId);
|
||||||
|
|
||||||
List<CompanySymbolNewRequest> companySymbolNewRequests = issuerCompanySymbolsList.stream().
|
List<CompanySymbolNewRequest> companySymbolNewRequests = issuerCompanySymbolsList.stream().
|
||||||
map(securityMkrRequestAdapter::toCompanySymbolRequest).toList();
|
map(issuerCompanyRequestAdapter::toCompanySymbolRequest).toList();
|
||||||
List<ContactNewRequest> contactNewRequests = issuerContactList.stream().
|
List<ContactNewRequest> contactNewRequests = issuerContactList.stream().
|
||||||
map(securityMkrRequestAdapter::toContactRequest).toList();
|
map(issuerCompanyRequestAdapter::toContactRequest).toList();
|
||||||
|
|
||||||
SecurityMkrGatewayRequest securityMkrGatewayRequest = new SecurityMkrGatewayRequest();
|
SecurityMkrGatewayRequest securityMkrGatewayRequest = new SecurityMkrGatewayRequest();
|
||||||
securityMkrGatewayRequest.setCompanyNewRequest(securityMkrRequestAdapter.toCompanyRequest(issuerCompany));
|
securityMkrGatewayRequest.setCompanyNewRequest(issuerCompanyRequestAdapter.toCompanyRequest(issuerCompany));
|
||||||
securityMkrGatewayRequest.setCompanySymbolNewRequests(companySymbolNewRequests);
|
securityMkrGatewayRequest.setCompanySymbolNewRequests(companySymbolNewRequests);
|
||||||
securityMkrGatewayRequest.setContactNewRequests(contactNewRequests);
|
securityMkrGatewayRequest.setContactNewRequests(contactNewRequests);
|
||||||
request.getIssuerCompanyInfoList()
|
request.getIssuerCompanyInfoList()
|
||||||
.stream()
|
.stream()
|
||||||
.filter(info -> info.getCompanyId().equals(companyId)).findFirst()
|
.filter(info -> info.getCompanyId().equals(companyId)).findFirst()
|
||||||
.map(securityMkrRequestAdapter::toCompanyInfoRequest)
|
.map(issuerCompanyRequestAdapter::toCompanyInfoRequest)
|
||||||
.ifPresent(securityMkrGatewayRequest::setCompanyInfoUpdateRequest);
|
.ifPresent(securityMkrGatewayRequest::setCompanyInfoUpdateRequest);
|
||||||
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.DESTINATION_ISSUER_COMPANY_GATEWAY_REQUEST, securityMkrGatewayRequest);
|
kafkaSender.sendRequestToQueue(Consts.DESTINATION_ISSUER_COMPANY_GATEWAY_REQUEST, securityMkrGatewayRequest);
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,41 @@
|
||||||
|
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.listing.mkr.ExchangeInstrument;
|
||||||
|
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.MMListingsRequest;
|
||||||
|
import ru.spcex.clearing.gatewayapi.service.adapter.MoneyMarketSecurityRequestAdapter;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.UUID;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class ListingMMProcessor {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
|
private final MoneyMarketSecurityRequestAdapter moneyMarketSecurityRequestAdapter;
|
||||||
|
private final KafkaSender kafkaSender;
|
||||||
|
|
||||||
|
public ListingMMProcessor(MoneyMarketSecurityRequestAdapter moneyMarketSecurityRequestAdapter,
|
||||||
|
KafkaSender kafkaSender) {
|
||||||
|
this.moneyMarketSecurityRequestAdapter = moneyMarketSecurityRequestAdapter;
|
||||||
|
this.kafkaSender = kafkaSender;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void process(MMListingsRequest request) {
|
||||||
|
List<ExchangeInstrument> exchangeInstrumentList = request.getExchangeInstrumentList();
|
||||||
|
for (ExchangeInstrument exchangeInstrument : exchangeInstrumentList) {
|
||||||
|
UUID exchangeInstrumentId = exchangeInstrument.getId();
|
||||||
|
log.debug("Process exchangeInstrumentId : {}", exchangeInstrumentId);
|
||||||
|
|
||||||
|
MoneyMarketSecurityNewRequest moneyMarketSecurityNewRequest = moneyMarketSecurityRequestAdapter
|
||||||
|
.toMoneyMarketSecurityNewRequest(exchangeInstrument);
|
||||||
|
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.DESTINATION_GATEWAY_MONEY_MARKET_SECURITY, moneyMarketSecurityNewRequest);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -5,7 +5,7 @@ import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.WithSecurityId;
|
import ru.spcex.clearing.gatewayapi.controller.request.WithSecurityId;
|
||||||
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.*;
|
import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.*;
|
||||||
import ru.spcex.clearing.gatewayapi.service.adapter.SecurityFondRequestAdapter;
|
import ru.spcex.clearing.gatewayapi.service.adapter.SecurityRequestAdapter;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
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.EquitySecurityGatewayRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.company.FixedIncomeGatewayRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.company.FixedIncomeGatewayRequest;
|
||||||
|
|
@ -22,10 +22,10 @@ import java.util.stream.Collectors;
|
||||||
public class SecurityService {
|
public class SecurityService {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
private final SecurityFondRequestAdapter securityRequestAdapter;
|
private final SecurityRequestAdapter securityRequestAdapter;
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
|
|
||||||
public SecurityService(SecurityFondRequestAdapter securityRequestAdapter,
|
public SecurityService(SecurityRequestAdapter securityRequestAdapter,
|
||||||
KafkaSender kafkaSender) {
|
KafkaSender kafkaSender) {
|
||||||
this.securityRequestAdapter = securityRequestAdapter;
|
this.securityRequestAdapter = securityRequestAdapter;
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
|
|
|
||||||
|
|
@ -12,7 +12,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.company.ContactNewRequest
|
||||||
import ru.spcex.platform.enumeration.CompanySymbol;
|
import ru.spcex.platform.enumeration.CompanySymbol;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class SecurityMkrRequestAdapter {
|
public class IssuerCompanyRequestAdapter {
|
||||||
|
|
||||||
public CompanyNewRequest toCompanyRequest(IssuerCompany issuerCompany) {
|
public CompanyNewRequest toCompanyRequest(IssuerCompany issuerCompany) {
|
||||||
CompanyNewRequest companyNewRequest = new CompanyNewRequest();
|
CompanyNewRequest companyNewRequest = new CompanyNewRequest();
|
||||||
|
|
@ -0,0 +1,23 @@
|
||||||
|
package ru.spcex.clearing.gatewayapi.service.adapter;
|
||||||
|
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.spcex.clearing.gatewayapi.controller.request.listing.mkr.ExchangeInstrument;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
|
||||||
|
import ru.spcex.platform.enumeration.InstrumentType;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class MoneyMarketSecurityRequestAdapter {
|
||||||
|
|
||||||
|
public MoneyMarketSecurityNewRequest toMoneyMarketSecurityNewRequest(ExchangeInstrument exchangeInstrument) {
|
||||||
|
MoneyMarketSecurityNewRequest moneyMarketSecurityNewRequest = new MoneyMarketSecurityNewRequest();
|
||||||
|
moneyMarketSecurityNewRequest.setNominalCurrency(exchangeInstrument.getSettlementCurrencyLetterCode());
|
||||||
|
moneyMarketSecurityNewRequest.setInstrumentType(InstrumentType.RATE.getKey());
|
||||||
|
moneyMarketSecurityNewRequest.setShortName(exchangeInstrument.getName());
|
||||||
|
moneyMarketSecurityNewRequest.setSecuritySymbol(exchangeInstrument.getSecuritySymbol());
|
||||||
|
moneyMarketSecurityNewRequest.setLotSize(exchangeInstrument.getLot());
|
||||||
|
moneyMarketSecurityNewRequest.setTermType(exchangeInstrument.getBankDepositAgreementType());
|
||||||
|
moneyMarketSecurityNewRequest.setIssuerId(exchangeInstrument.getIssuerId());
|
||||||
|
moneyMarketSecurityNewRequest.setWorkflowStatus(exchangeInstrument.getWorkflowStatus());
|
||||||
|
return moneyMarketSecurityNewRequest;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -8,7 +8,7 @@ import ru.spcex.clearing.gatewayapi.controller.request.listing.fond.Nominal;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.*;
|
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.*;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class SecurityFondRequestAdapter {
|
public class SecurityRequestAdapter {
|
||||||
|
|
||||||
public FixedIncomeSecurityNewRequest toFixedIncomeSecurityRequest(FondSecurity fondSecurity) {
|
public FixedIncomeSecurityNewRequest toFixedIncomeSecurityRequest(FondSecurity fondSecurity) {
|
||||||
FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest = new FixedIncomeSecurityNewRequest();
|
FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest = new FixedIncomeSecurityNewRequest();
|
||||||
|
|
@ -2,6 +2,7 @@ package ru.spcex.clearing.platform.messaging.domain;
|
||||||
|
|
||||||
public interface Consts {
|
public interface Consts {
|
||||||
String DESTINATION_MONEY_MARKET_SECURITY_NEW = "money-market-security-new";
|
String DESTINATION_MONEY_MARKET_SECURITY_NEW = "money-market-security-new";
|
||||||
|
String DESTINATION_GATEWAY_MONEY_MARKET_SECURITY = "money-gateway-market-security";
|
||||||
String DESTINATION_MONEY_MARKET_SECURITY_UPDATE = "money-market-security-update";
|
String DESTINATION_MONEY_MARKET_SECURITY_UPDATE = "money-market-security-update";
|
||||||
String DESTINATION_MONEY_MARKET_SECURITY_DELETE = "money-market-security-delete";
|
String DESTINATION_MONEY_MARKET_SECURITY_DELETE = "money-market-security-delete";
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -5,9 +5,9 @@ import java.math.BigDecimal;
|
||||||
public class AssetOperationRequest {
|
public class AssetOperationRequest {
|
||||||
private BigDecimal amount;
|
private BigDecimal amount;
|
||||||
private BigDecimal quantity;
|
private BigDecimal quantity;
|
||||||
private String asset;
|
private String securitySymbol;
|
||||||
private String code;
|
private String code;
|
||||||
private String firmId;
|
private String tradingCode;
|
||||||
|
|
||||||
public BigDecimal getAmount() {
|
public BigDecimal getAmount() {
|
||||||
return amount;
|
return amount;
|
||||||
|
|
@ -25,12 +25,12 @@ public class AssetOperationRequest {
|
||||||
this.quantity = quantity;
|
this.quantity = quantity;
|
||||||
}
|
}
|
||||||
|
|
||||||
public String getAsset() {
|
public String getSecuritySymbol() {
|
||||||
return asset;
|
return securitySymbol;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setAsset(String asset) {
|
public void setSecuritySymbol(String securitySymbol) {
|
||||||
this.asset = asset;
|
this.securitySymbol = securitySymbol;
|
||||||
}
|
}
|
||||||
|
|
||||||
public String getCode() {
|
public String getCode() {
|
||||||
|
|
@ -41,11 +41,11 @@ public class AssetOperationRequest {
|
||||||
this.code = code;
|
this.code = code;
|
||||||
}
|
}
|
||||||
|
|
||||||
public String getFirmId() {
|
public String getTradingCode() {
|
||||||
return firmId;
|
return tradingCode;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setFirmId(String firmId) {
|
public void setTradingCode(String tradingCode) {
|
||||||
this.firmId = firmId;
|
this.tradingCode = tradingCode;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue