This commit is contained in:
parent
10cf6bc3cb
commit
270a27187c
6 changed files with 124 additions and 24 deletions
|
|
@ -24,6 +24,8 @@ import ru.spcex.clearing.backendapi.security.KeycloakUtils;
|
|||
import ru.spcex.clearing.backendapi.service.IOperator;
|
||||
import ru.spcex.clearing.backendapi.service.IStateLoader;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
import ru.spcex.platform.enumeration.Task;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||
|
|
@ -67,7 +69,7 @@ public class LauncherController extends AbstractQueueController {
|
|||
@RequestMapping(value = "/{task-code}", method = RequestMethod.POST, produces = MediaType.APPLICATION_JSON_VALUE)
|
||||
@ResponseBody
|
||||
public CudResponse add(@ApiParam(value = "Код задания из taskDictionary", required = true, example = "ABLK")
|
||||
@PathVariable("task-code") String dictionaryName) throws ExecutionException, InterruptedException {
|
||||
@PathVariable("task-code") String dictionaryName) throws ExecutionException, InterruptedException {
|
||||
AbstractDictionary taskEnum = taskDictionary.getSingleObjectByFieldValues(Map.of("code", dictionaryName));
|
||||
if (taskEnum == null) {
|
||||
throw new NotFound404Exception("task dictionary element with code '" + dictionaryName + "'");
|
||||
|
|
@ -81,11 +83,37 @@ public class LauncherController extends AbstractQueueController {
|
|||
}
|
||||
launcherCommand.setTask(dictionaryName);
|
||||
launcherCommand.setUserId(user.getId());
|
||||
// пока здесь, это требуется для сохранения истории
|
||||
saveLauncher(dictionaryName, user.getId());
|
||||
//топики ограничиваются наличием в taskDictionary
|
||||
//подписываются на разные топики в разных модулях, см. ru.spcex.platform.enumeration.Task#topic
|
||||
return processRequest("launcher-" + dictionaryName, launcherCommand);
|
||||
return processRequest(Consts.LAUNCHER_NEW, launcherCommand);
|
||||
}
|
||||
|
||||
@ApiOperation(value = "create specific launcher.")
|
||||
@ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = CudResponse.class), @ApiResponse(code = 400, message = "Ошибка валидации", response = BasicSpcexResponse.class)})
|
||||
@RequestMapping(value = "/specific/{task-code}", method = RequestMethod.POST, produces = MediaType.APPLICATION_JSON_VALUE)
|
||||
@ResponseBody
|
||||
public CudResponse addSpecific(@ApiParam(value = "Код задания из taskDictionary", required = true, example = "ABLK")
|
||||
@PathVariable("task-code") String dictionaryName,
|
||||
@RequestBody LauncherNew launcherNew) throws ExecutionException, InterruptedException {
|
||||
AbstractDictionary taskEnum = taskDictionary.getSingleObjectByFieldValues(Map.of("code", dictionaryName));
|
||||
if (taskEnum == null) {
|
||||
throw new NotFound404Exception("task dictionary element with code '" + dictionaryName + "'");
|
||||
}
|
||||
if (!Task.startOfClearing.getKey().equals(taskEnum.getCode())) {
|
||||
throw new IllegalStateException(String.format("Task %s not support request with body", taskEnum.getCode()));
|
||||
}
|
||||
LauncherNew launcherCommand = new LauncherNew();
|
||||
Authentication authentication = SecurityContextHolder.getContext().getAuthentication();
|
||||
String username = KeycloakUtils.getUserNameFromAuthentication(authentication);
|
||||
User user = userImdg.getSingleObjectByFieldValues(Map.of("identifier", username));
|
||||
if (user == null) {
|
||||
throw new IllegalStateException("cannot obtain userId from logged in user " + username);
|
||||
}
|
||||
launcherCommand.setTask(dictionaryName);
|
||||
launcherCommand.setUserId(user.getId());
|
||||
launcherCommand.setCompanyId(launcherNew.getCompanyId());
|
||||
launcherCommand.setSecurityId(launcherNew.getSecurityId());
|
||||
return processRequest(Consts.LAUNCHER_NEW, launcherCommand);
|
||||
}
|
||||
|
||||
private void saveLauncher(String taskCode, Long userId) {
|
||||
|
|
|
|||
|
|
@ -19,6 +19,12 @@ public class LauncherNew implements IAction<Object> {
|
|||
@ApiModelProperty(value = "Идентификатор единоличного исполнительного органа", example = "ABCD")
|
||||
@JsonProperty
|
||||
private String task;
|
||||
@ApiModelProperty(value = "Идентификатор инициатора", example = "1000")
|
||||
@JsonProperty
|
||||
private Long companyId;
|
||||
@ApiModelProperty(value = "Идентификатор инструмента", example = "1000")
|
||||
@JsonProperty
|
||||
private Long securityId;
|
||||
|
||||
@JsonIgnore
|
||||
private Long userId;
|
||||
|
|
@ -32,6 +38,8 @@ public class LauncherNew implements IAction<Object> {
|
|||
} else {
|
||||
LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest();
|
||||
taskRunnerCommandRequest.setTaskName(task);
|
||||
taskRunnerCommandRequest.setCompanyId(companyId);
|
||||
taskRunnerCommandRequest.setSecurityId(securityId);
|
||||
taskRunnerCommandRequest.setUserId(userId);
|
||||
return taskRunnerCommandRequest;
|
||||
}
|
||||
|
|
@ -65,4 +73,20 @@ public class LauncherNew implements IAction<Object> {
|
|||
public void setUserId(Long userId) {
|
||||
this.userId = userId;
|
||||
}
|
||||
|
||||
public Long getCompanyId() {
|
||||
return companyId;
|
||||
}
|
||||
|
||||
public void setCompanyId(Long companyId) {
|
||||
this.companyId = companyId;
|
||||
}
|
||||
|
||||
public Long getSecurityId() {
|
||||
return securityId;
|
||||
}
|
||||
|
||||
public void setSecurityId(Long securityId) {
|
||||
this.securityId = securityId;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,6 +1,5 @@
|
|||
package ru.spcex.clearing.service;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
||||
import org.apache.kafka.clients.producer.RecordMetadata;
|
||||
|
|
@ -10,15 +9,11 @@ import org.springframework.beans.factory.annotation.Autowired;
|
|||
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.util.StringUtils;
|
||||
import ru.clearing.classes.statics.data.account.Account;
|
||||
import ru.clearing.classes.statics.data.account.AccountBalance;
|
||||
import ru.clearing.classes.statics.data.clearing.VerificationResult;
|
||||
import ru.clearing.classes.statics.data.company.Company;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
||||
import ru.clearing.classes.statics.data.misc.Listing;
|
||||
import ru.clearing.classes.statics.data.misc.STrade;
|
||||
import ru.clearing.classes.statics.data.sdf.SDf01;
|
||||
import ru.clearing.classes.statics.data.security.Security;
|
||||
import ru.spcex.clearing.error.ClearingError;
|
||||
import ru.spcex.clearing.error.ClearingException;
|
||||
|
|
@ -26,12 +21,10 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
|||
import ru.spcex.clearing.platform.messaging.domain.ActionType;
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.CoveredDealRegisterNewRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityNewRequest;
|
||||
import ru.spcex.platform.classes.base.interfaces.WithId;
|
||||
import ru.spcex.platform.enumeration.AccountType;
|
||||
import ru.spcex.platform.enumeration.Allowed;
|
||||
import ru.spcex.platform.enumeration.Market;
|
||||
import ru.spcex.platform.enumeration.ResultStatuses;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgId;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
|
@ -40,17 +33,13 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
|||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||
import ru.spcex.platform.utils.time.TimeUtil;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.CoveredDealRegisterNewRequest;
|
||||
|
||||
|
||||
import java.math.BigDecimal;
|
||||
import java.math.RoundingMode;
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
|
|
@ -96,7 +85,7 @@ public class ExecutionDepositComponent {
|
|||
/**
|
||||
* Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?"
|
||||
*/
|
||||
@Scheduled(cron = "${clearing-service.scheduler.check-s-trade}")
|
||||
@Scheduled(cron = "${clearing-service.scheduler.check-s-trade}")
|
||||
public void resetTradingDay() {
|
||||
Instant today = TimeUtil.localDateToInstant(LocalDate.now());
|
||||
if (tradingDay == null || !tradingDay.equals(today)) {
|
||||
|
|
@ -109,7 +98,7 @@ public class ExecutionDepositComponent {
|
|||
public void processNewTS() {
|
||||
Long tradeNum = -1L; // todo уточнить как он обновляется
|
||||
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
|
||||
ImdgPredicate sql = pb.and(pb.greater("tradeNum",tradeNum), pb.greatEqual("tradeDateTime", tradingDay));
|
||||
ImdgPredicate sql = pb.and(pb.greater("tradeNum", tradeNum), pb.greatEqual("tradeDateTime", tradingDay));
|
||||
Collection<STrade> sTrades = sTradeImdg.getCollectionObjectsByPredicate(sql);
|
||||
log.info("Found {} new s_trade with trade_num>{}", sTrades.size(), tradeNum);
|
||||
|
||||
|
|
@ -145,7 +134,7 @@ public class ExecutionDepositComponent {
|
|||
|
||||
Long generationId = idGenerator.nextId();
|
||||
log.info("generationId = {}", generationId);
|
||||
for (STrade trade:sTrades) {
|
||||
for (STrade trade : sTrades) {
|
||||
log.trace("Check s_trade[{}].tradeNum={}", trade.getId(), trade.getTradeNum());
|
||||
Collection<ExecutionDeposit> existsEDeposit = executionDepositImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
"exchangeExecutionId", trade.getTradeNum(),
|
||||
|
|
@ -170,13 +159,13 @@ public class ExecutionDepositComponent {
|
|||
}
|
||||
|
||||
} else {
|
||||
long[] idToLong = existsEDeposit.stream().mapToLong(ed-> ed.getId()).toArray();
|
||||
long[] idToLong = existsEDeposit.stream().mapToLong(ed -> ed.getId()).toArray();
|
||||
log.warn("S_TRADE[{}] already has executionDeposit: {}", trade.getId(), Arrays.toString(idToLong));
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Long newMaxTradeNum = sTrades.stream().mapToLong(STrade::getTradeNum).max().orElseGet(()-> tradeNum);
|
||||
Long newMaxTradeNum = sTrades.stream().mapToLong(STrade::getTradeNum).max().orElseGet(() -> tradeNum);
|
||||
log.debug("Next tradeNum is {}", newMaxTradeNum);
|
||||
}
|
||||
|
||||
|
|
@ -203,13 +192,14 @@ public class ExecutionDepositComponent {
|
|||
|
||||
/**
|
||||
* в очередь kafka для модуля securities-service сообщение о добавлении инструмента с параметром securitySymbol=s_trade.sec_code
|
||||
*
|
||||
* @param newSymbolRequest
|
||||
*/
|
||||
protected void createNewSecurities(Collection<String> newSymbolRequest) {
|
||||
final String destination = Consts.DESTINATION_MONEY_MARKET_SECURITY_NEW;
|
||||
List<String> symbolRequests = new ArrayList<>(newSymbolRequest); // чтобы в случае ошибки отобразить номер в логе
|
||||
List<Future<RecordMetadata>> sendAll = new ArrayList<>(symbolRequests.size());
|
||||
for (String newSymbol: symbolRequests) {
|
||||
for (String newSymbol : symbolRequests) {
|
||||
if (newSymbol == null || newSymbol.isEmpty()) {
|
||||
log.warn("Empty SecuritySumbol");
|
||||
} else {
|
||||
|
|
@ -228,7 +218,7 @@ public class ExecutionDepositComponent {
|
|||
}
|
||||
}
|
||||
int i = 0;
|
||||
for (Future<RecordMetadata> future: sendAll) {
|
||||
for (Future<RecordMetadata> future : sendAll) {
|
||||
try {
|
||||
future.get(); // get exception
|
||||
} catch (InterruptedException e) {
|
||||
|
|
|
|||
|
|
@ -0,0 +1,33 @@
|
|||
package ru.spcex.clearing.service;
|
||||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.platform.enumeration.Task;
|
||||
|
||||
@Service
|
||||
public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
|
||||
private final ClearingService clearingService;
|
||||
@Autowired
|
||||
public LauncherCommandReceiver(Consumer<String, Object> kafkaQueue,
|
||||
ClearingService clearingService) {
|
||||
super(kafkaQueue);
|
||||
this.clearingService = clearingService;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() {
|
||||
callback(LauncherCommandRequest.class)
|
||||
.setConsumer(action -> clearingService.executeVerification())
|
||||
.forDestination(Task.createOrderConfirm.topic(), callbacks::put); // CORC
|
||||
|
||||
init();
|
||||
}
|
||||
}
|
||||
|
|
@ -2,10 +2,12 @@ package ru.spcex.clearing.scheduler.service;
|
|||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.scheduler.Launcher;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
|
|
@ -20,14 +22,17 @@ import java.time.Instant;
|
|||
|
||||
import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey;
|
||||
|
||||
@Service
|
||||
public class LauncherService extends QueueConsumer implements InitializingBean {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final Imdg<Launcher> launcherMap;
|
||||
private final Producer<String, Object> kafkaProducer;
|
||||
|
||||
@Autowired
|
||||
public LauncherService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
|
||||
ImdgProvider imdgProvider) {
|
||||
super(kafkaQueue, kafkaProducer);
|
||||
this.kafkaProducer = kafkaProducer;
|
||||
this.launcherMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Launcher, Launcher.class);
|
||||
}
|
||||
|
||||
|
|
@ -53,6 +58,7 @@ public class LauncherService extends QueueConsumer implements InitializingBean {
|
|||
launcher.setUpdated(created);
|
||||
launcherMap.insert(launcher);
|
||||
log.debug("successfully processed, new id {}", launcher.getId());
|
||||
kafkaProducer.send(new ProducerRecord<>("launcher-" + launcher.getTask(), req));
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,9 +6,12 @@ public class LauncherCommandRequest {
|
|||
|
||||
@JsonProperty
|
||||
private Long userId;
|
||||
|
||||
@JsonProperty
|
||||
private String taskName;
|
||||
@JsonProperty
|
||||
private Long companyId;
|
||||
@JsonProperty
|
||||
private Long securityId;
|
||||
|
||||
|
||||
public Long getUserId() {
|
||||
|
|
@ -26,4 +29,20 @@ public class LauncherCommandRequest {
|
|||
public void setTaskName(String taskName) {
|
||||
this.taskName = taskName;
|
||||
}
|
||||
|
||||
public Long getCompanyId() {
|
||||
return companyId;
|
||||
}
|
||||
|
||||
public void setCompanyId(Long companyId) {
|
||||
this.companyId = companyId;
|
||||
}
|
||||
|
||||
public Long getSecurityId() {
|
||||
return securityId;
|
||||
}
|
||||
|
||||
public void setSecurityId(Long securityId) {
|
||||
this.securityId = securityId;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue