This commit is contained in:
parent
224179a9e1
commit
dcc0a3c7ce
1 changed files with 41 additions and 13 deletions
|
|
@ -6,6 +6,7 @@ import java.util.Collection;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Objects;
|
import java.util.Objects;
|
||||||
|
import java.util.Optional;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
|
|
@ -40,6 +41,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.gateway.Tkr;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.reports.NotificationRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.reports.NotificationRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
|
|
@ -49,6 +51,7 @@ import ru.spcex.clearing.util.services.RequestHelper;
|
||||||
import ru.spcex.clearing.validation.common.ValidationHelper;
|
import ru.spcex.clearing.validation.common.ValidationHelper;
|
||||||
import ru.spcex.platform.enumeration.CompanySymbol;
|
import ru.spcex.platform.enumeration.CompanySymbol;
|
||||||
import ru.spcex.platform.enumeration.ServiceStatus;
|
import ru.spcex.platform.enumeration.ServiceStatus;
|
||||||
|
import static ru.spcex.platform.enumeration.Task.makeFiles_MTCR;
|
||||||
import ru.spcex.platform.enumeration.TradingClearingRegistryPurpose;
|
import ru.spcex.platform.enumeration.TradingClearingRegistryPurpose;
|
||||||
import ru.spcex.platform.enumeration.TradingClearingRegistryType;
|
import ru.spcex.platform.enumeration.TradingClearingRegistryType;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
|
@ -142,6 +145,9 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
callback(TkrAccountsGatewayRequest.class)
|
callback(TkrAccountsGatewayRequest.class)
|
||||||
.setFunction(this::tradingClearingRegistryCheck)
|
.setFunction(this::tradingClearingRegistryCheck)
|
||||||
.forDestination(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_CHECK, callbacks::put);
|
.forDestination(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_CHECK, callbacks::put);
|
||||||
|
callback(LauncherCommandRequest.class)
|
||||||
|
.setConsumer(this::sendAllTkrToGateway)
|
||||||
|
.forDestination(makeFiles_MTCR.topic(), callbacks::put);
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -229,8 +235,10 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
targetExistTradingClearingRegistry.setUpdated(Instant.now());
|
targetExistTradingClearingRegistry.setUpdated(Instant.now());
|
||||||
tradingClearingRegistryImdg.update(targetExistTradingClearingRegistry);
|
tradingClearingRegistryImdg.update(targetExistTradingClearingRegistry);
|
||||||
log.debug("successfully processed, id {}. Updated exist TradingClearingRegistry.id={}", id, targetExistTradingClearingRegistry.getId());
|
log.debug("successfully processed, id {}. Updated exist TradingClearingRegistry.id={}", id, targetExistTradingClearingRegistry.getId());
|
||||||
SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry);
|
Optional<Tkr> tkr = createRequestToGateway(tradingClearingRegistry);
|
||||||
if (!sendTkrRequest.getTkrs().isEmpty()) {
|
if (tkr.isPresent()) {
|
||||||
|
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||||
|
sendTkrRequest.getTkrs().add(tkr.get());
|
||||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||||
}
|
}
|
||||||
return null;
|
return null;
|
||||||
|
|
@ -367,8 +375,10 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
|
|
||||||
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
|
tradingClearingRegistryImdg.insert(tradingClearingRegistry);
|
||||||
log.info("New TCR.id={} was created.", tradingClearingRegistry.getId());
|
log.info("New TCR.id={} was created.", tradingClearingRegistry.getId());
|
||||||
SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry);
|
Optional<Tkr> tkr = createRequestToGateway(tradingClearingRegistry);
|
||||||
if (!sendTkrRequest.getTkrs().isEmpty()) {
|
if (tkr.isPresent()) {
|
||||||
|
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||||
|
sendTkrRequest.getTkrs().add(tkr.get());
|
||||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -462,8 +472,10 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
tradingClearingRegistry.setUpdated(now);
|
tradingClearingRegistry.setUpdated(now);
|
||||||
tradingClearingRegistry.setStatus(req.getStatus());
|
tradingClearingRegistry.setStatus(req.getStatus());
|
||||||
tradingClearingRegistryImdg.update(tradingClearingRegistry);
|
tradingClearingRegistryImdg.update(tradingClearingRegistry);
|
||||||
SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry);
|
Optional<Tkr> tkr = createRequestToGateway(tradingClearingRegistry);
|
||||||
if (!sendTkrRequest.getTkrs().isEmpty()) {
|
if (tkr.isPresent()) {
|
||||||
|
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||||
|
sendTkrRequest.getTkrs().add(tkr.get());
|
||||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||||
}
|
}
|
||||||
log.debug("Update TCR.id={}.", tradingClearingRegistry.getId());
|
log.debug("Update TCR.id={}.", tradingClearingRegistry.getId());
|
||||||
|
|
@ -503,18 +515,31 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
TradingClearingRegistry tradingClearingRegistry = tradingClearingRegistryImdg.getSingleObjectByFieldValues(
|
TradingClearingRegistry tradingClearingRegistry = tradingClearingRegistryImdg.getSingleObjectByFieldValues(
|
||||||
Map.of("code", accountForCheck.getTkrCode())
|
Map.of("code", accountForCheck.getTkrCode())
|
||||||
);
|
);
|
||||||
SendTkrRequest sendTkrRequest = createRequestToGateway(tradingClearingRegistry);
|
Optional<Tkr> tkr = createRequestToGateway(tradingClearingRegistry);
|
||||||
sendTkrRequest.setRequestId(req.getRequestId());
|
if (tkr.isPresent()) {
|
||||||
if (!sendTkrRequest.getTkrs().isEmpty()) {
|
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||||
|
sendTkrRequest.getTkrs().add(tkr.get());
|
||||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
private SendTkrRequest createRequestToGateway(TradingClearingRegistry tradingClearingRegistry) {
|
public RequestInfoUpdate sendAllTkrToGateway(BaseRequest<?> gatewayRequest) {
|
||||||
|
Collection<TradingClearingRegistry> tkrs = tradingClearingRegistryImdg.getAllValues();
|
||||||
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||||
|
for (TradingClearingRegistry tkr : tkrs) {
|
||||||
|
log.debug("Create TKR request with id: {}", tkr.getId());
|
||||||
|
Optional<Tkr> tkrRequest = createRequestToGateway(tkr);
|
||||||
|
tkrRequest.ifPresent(value -> sendTkrRequest.getTkrs().add(value));
|
||||||
|
}
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
private Optional<Tkr> createRequestToGateway(TradingClearingRegistry tradingClearingRegistry) {
|
||||||
if (tradingClearingRegistry != null) {
|
if (tradingClearingRegistry != null) {
|
||||||
|
Tkr tkr = new Tkr();
|
||||||
log.debug("Found TKR with id: {} and code: {}", tradingClearingRegistry.getId(), tradingClearingRegistry.getCode());
|
log.debug("Found TKR with id: {} and code: {}", tradingClearingRegistry.getId(), tradingClearingRegistry.getCode());
|
||||||
Company company = companyImdg.getSingleObjectByID(tradingClearingRegistry.getCompanyId());
|
Company company = companyImdg.getSingleObjectByID(tradingClearingRegistry.getCompanyId());
|
||||||
CompanySymbols companySymbols = companySymbolsImdg.getSingleObjectByFieldValues(
|
CompanySymbols companySymbols = companySymbolsImdg.getSingleObjectByFieldValues(
|
||||||
|
|
@ -523,10 +548,13 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
"companyId", company.getId()
|
"companyId", company.getId()
|
||||||
)
|
)
|
||||||
);
|
);
|
||||||
|
if (companySymbols == null){
|
||||||
|
log.warn("CompanySymbols is null, searched by company.id: {}, skip this tkr", company.getId());
|
||||||
|
return Optional.empty();
|
||||||
|
}
|
||||||
ClientCode clientCode = clientCodeImdg.getFirstObjectByFieldValues(
|
ClientCode clientCode = clientCodeImdg.getFirstObjectByFieldValues(
|
||||||
Map.of("tradingClearingRegistryId", tradingClearingRegistry.getId())
|
Map.of("tradingClearingRegistryId", tradingClearingRegistry.getId())
|
||||||
);
|
);
|
||||||
Tkr tkr = new Tkr();
|
|
||||||
tkr.setCompanyId(companySymbols.getCompanySymbolValue());
|
tkr.setCompanyId(companySymbols.getCompanySymbolValue());
|
||||||
tkr.setTradingCode(Long.valueOf(company.getTradingCode()));
|
tkr.setTradingCode(Long.valueOf(company.getTradingCode()));
|
||||||
tkr.setTkrCode(tradingClearingRegistry.getCode());
|
tkr.setTkrCode(tradingClearingRegistry.getCode());
|
||||||
|
|
@ -575,9 +603,9 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
||||||
tkr.getMoneyAccounts().add(moneyAccountMsg);
|
tkr.getMoneyAccounts().add(moneyAccountMsg);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
sendTkrRequest.getTkrs().add(tkr);
|
return Optional.of(tkr);
|
||||||
}
|
}
|
||||||
return sendTkrRequest;
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue