http://jira.mfd.msk:8088/browse/CLS-958 [1 - send orders by fix protocol]
This commit is contained in:
parent
63ec986491
commit
4ad5a13710
3 changed files with 74 additions and 198 deletions
|
|
@ -114,7 +114,7 @@ public class TaskListener extends QueueConsumer implements InitializingBean {
|
||||||
}
|
}
|
||||||
|
|
||||||
private void exportToTriAndSendToFix(BaseRequest<LauncherCommandRequest> req) {
|
private void exportToTriAndSendToFix(BaseRequest<LauncherCommandRequest> req) {
|
||||||
|
log.info("SDOR export and send command to fix, request={}", req);
|
||||||
try {
|
try {
|
||||||
Collection<OrderCurrency> orders = imdgQueryService.getClccOrderCurrencyIds();
|
Collection<OrderCurrency> orders = imdgQueryService.getClccOrderCurrencyIds();
|
||||||
String content = orderCurrencyTriExportService.exportToTri(orders);
|
String content = orderCurrencyTriExportService.exportToTri(orders);
|
||||||
|
|
|
||||||
|
|
@ -1,197 +0,0 @@
|
||||||
package ru.spcex.clearing.fix.service;
|
|
||||||
|
|
||||||
import org.apache.commons.lang3.StringUtils;
|
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
|
||||||
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.company.Company;
|
|
||||||
import ru.clearing.classes.statics.data.company.CompanySymbols;
|
|
||||||
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.cud.company.CompanyNewRequest;
|
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
|
||||||
import ru.spcex.clearing.util.security.UserRoleVerification;
|
|
||||||
import ru.spcex.clearing.util.services.RequestHelper;
|
|
||||||
import ru.spcex.platform.enumeration.CompanySymbol;
|
|
||||||
import ru.spcex.platform.enumeration.UserRole;
|
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgId;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgTransaction;
|
|
||||||
import ru.spcex.platform.utils.error.ValidationException;
|
|
||||||
|
|
||||||
import java.time.Instant;
|
|
||||||
import java.util.Objects;
|
|
||||||
|
|
||||||
@Service
|
|
||||||
public class FixService extends QueueConsumer implements InitializingBean {
|
|
||||||
|
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
|
||||||
private final RequestHelper requestHelper;
|
|
||||||
private final ImdgProvider imdgProvider;
|
|
||||||
private final ImdgId idSequence;
|
|
||||||
private final Imdg<Company> companyIMap;
|
|
||||||
protected UserRoleVerification userRoleVerification;
|
|
||||||
// protected CompanySymbolService companySymbolService;
|
|
||||||
// protected AccountNotificationHelper accountNotification;
|
|
||||||
// protected RelationService relationService;
|
|
||||||
// private final ValidationHelper validationHelper;
|
|
||||||
// private final Function<CompanyNewRequest, IValidator> companyNewRequestValidator;
|
|
||||||
// private final Function<CompanyNewRequest, IValidator> companyUpdateRequestValidator;
|
|
||||||
// private final Function<CommonDeleteRequest, IValidator> companyDeleteRequestValidator;
|
|
||||||
|
|
||||||
@Autowired
|
|
||||||
public FixService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
|
|
||||||
ImdgProvider imdgProvider,
|
|
||||||
RequestHelper requestHelper,
|
|
||||||
// ValidationHelper validationHelper,
|
|
||||||
UserRoleVerification userRoleVerification
|
|
||||||
// @Qualifier("companyNewRequestValidator")
|
|
||||||
// Function<CompanyNewRequest, IValidator> companyNewRequestValidator,
|
|
||||||
// @Qualifier("companyUpdateRequestValidator")
|
|
||||||
// Function<CompanyNewRequest, IValidator> companyUpdateRequestValidator,
|
|
||||||
// @Qualifier("CompanyDeleteRequestValidator")
|
|
||||||
// Function<CommonDeleteRequest, IValidator> companyDeleteRequestValidator,
|
|
||||||
// CompanySymbolService companySymbolService,
|
|
||||||
// AccountNotificationHelper accountNotification,
|
|
||||||
/*RelationService relationService*/) {
|
|
||||||
super(kafkaQueue, kafkaProducer);
|
|
||||||
this.imdgProvider = imdgProvider;
|
|
||||||
this.requestHelper = requestHelper.setLogger(log);
|
|
||||||
this.idSequence = imdgProvider.getImdgIdGenerator();
|
|
||||||
this.userRoleVerification = userRoleVerification;
|
|
||||||
this.userRoleVerification.setRoleForVerification(UserRole.Admin);
|
|
||||||
// this.validationHelper = validationHelper;
|
|
||||||
// this.companyNewRequestValidator = companyNewRequestValidator;
|
|
||||||
// this.companyUpdateRequestValidator = companyUpdateRequestValidator;
|
|
||||||
// this.companyDeleteRequestValidator = companyDeleteRequestValidator;
|
|
||||||
// this.companySymbolService = companySymbolService;
|
|
||||||
// this.companySymbolService.setCompanyService(this);
|
|
||||||
// this.accountNotification = accountNotification;
|
|
||||||
// this.relationService = relationService;
|
|
||||||
|
|
||||||
companyIMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public void afterPropertiesSet() {
|
|
||||||
// callback(CommonDeleteRequest.class)
|
|
||||||
// .setFunction(request -> requestHelper.requestFunction(this::deleteCompany, request))
|
|
||||||
// .forDestination(Consts.DESTINATION_COMPANY_DELETE, callbacks::put);
|
|
||||||
// callback(CommonDeleteRequest.class)
|
|
||||||
// .setFunction(request -> requestHelper.requestFunction(this::blockCompanyAfterDocument, request))
|
|
||||||
// .forDestination(Consts.DESTINATION_COMPANY_BLOCK, callbacks::put);
|
|
||||||
// callback(CompanyNewRequest.class)
|
|
||||||
// .setFunction(request -> requestHelper.requestFunction(this::processBaseRequest, request))
|
|
||||||
// .forDestination(Consts.DESTINATION_COMPANY_NEW, callbacks::put);
|
|
||||||
// callback(CompanyNewRequest.class)
|
|
||||||
// .setFunction(request -> requestHelper.requestFunction(this::processBaseRequest, request))
|
|
||||||
// .forDestination(Consts.DESTINATION_COMPANY_UPDATE, callbacks::put);
|
|
||||||
// callback(AccountTerminationRequest.class)
|
|
||||||
// .setFunction(request -> requestHelper.requestFunction(this::finishCompanyTermination, request))
|
|
||||||
// .forDestination(Consts.ACCOUNT_TERMINATION_STEP2, callbacks::put);
|
|
||||||
init();
|
|
||||||
}
|
|
||||||
|
|
||||||
private RequestInfoUpdate processBaseRequest(BaseRequest<?> request) throws ValidationException {
|
|
||||||
// Валидация
|
|
||||||
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(request);
|
|
||||||
if (requestInfoUpdate != null) return requestInfoUpdate;
|
|
||||||
ActionType actionType = request.getActionType();
|
|
||||||
switch (actionType) {
|
|
||||||
case NEW -> {
|
|
||||||
CompanyNewRequest companyNewRequest = (CompanyNewRequest) request.getRequestPayload();
|
|
||||||
// IValidator validator = companyNewRequestValidator.apply(companyNewRequest);
|
|
||||||
// Optional<EnumMessage> error = validator.tillFirstError();
|
|
||||||
// if (error.isPresent()) {
|
|
||||||
// throw new ValidationException(error.get());
|
|
||||||
// }
|
|
||||||
// create(companyNewRequest);
|
|
||||||
}
|
|
||||||
case UPDATE -> update((CompanyNewRequest) request.getRequestPayload(), true);
|
|
||||||
}
|
|
||||||
return requestInfoUpdate;
|
|
||||||
}
|
|
||||||
|
|
||||||
public synchronized void update(CompanyNewRequest updateRequest, boolean validationIsEnable) throws ValidationException {
|
|
||||||
// if (validationIsEnable) {
|
|
||||||
// IValidator validator = companyUpdateRequestValidator.apply(updateRequest);
|
|
||||||
// Optional<EnumMessage> error = validator.tillFirstError();
|
|
||||||
// if (error.isPresent()) {
|
|
||||||
// throw new ValidationException(error.get());
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
ImdgTransaction transaction = imdgProvider.newTransaction();
|
|
||||||
transaction.beginTransaction();
|
|
||||||
boolean txOk = false;
|
|
||||||
try {
|
|
||||||
Imdg<Company> companyMap = transaction.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
|
||||||
Company company = companyMap.getSingleObjectByID(updateRequest.getId());
|
|
||||||
// if (company == null) {
|
|
||||||
// log.trace("Company {} not found", updateRequest.getId());
|
|
||||||
// throw new ValidationException(new EnumMessage(FixErrors.CompanyNotFound));
|
|
||||||
// }
|
|
||||||
// if (!WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
|
|
||||||
// log.trace("Company {} not active: {}", company.getId(), company.getWorkflowStatus());
|
|
||||||
// return requestHelper.makeErrorResponse(companyUpdateRequestBaseRequest, CompanyErrors.CompanyDisabled, updateRequest.getId());
|
|
||||||
// }
|
|
||||||
|
|
||||||
company.setUpdated(Instant.now());
|
|
||||||
if (StringUtils.isNotEmpty(updateRequest.getShortName()))
|
|
||||||
company.setShortName(updateRequest.getShortName());
|
|
||||||
if (StringUtils.isNotEmpty(updateRequest.getFullName()))
|
|
||||||
company.setFullName(updateRequest.getFullName());
|
|
||||||
|
|
||||||
if (updateRequest.getCompanySymbol() != null || updateRequest.getCompanySymbolValue() != null) {
|
|
||||||
log.trace("Request field CompanySymbol, CompanySymbolValue ignore for update company request.");
|
|
||||||
}
|
|
||||||
|
|
||||||
String prevStatus = company.getWorkflowStatus();
|
|
||||||
if (updateRequest.getWorkflowStatus() != null) {
|
|
||||||
company.setWorkflowStatus(updateRequest.getWorkflowStatus());
|
|
||||||
if (!Objects.equals(prevStatus, company.getWorkflowStatus())) {
|
|
||||||
// relationService.onChangeWorkflowStatus(transaction, company, prevStatus, company.getWorkflowStatus());
|
|
||||||
} else {
|
|
||||||
log.trace("Status was not changed");
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
log.trace("Null new WorkflowStatus");
|
|
||||||
}
|
|
||||||
companyMap.update(company);
|
|
||||||
txOk = true;
|
|
||||||
} finally {
|
|
||||||
if (txOk)
|
|
||||||
transaction.commitTransaction();
|
|
||||||
else
|
|
||||||
transaction.rollbackTransaction();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
protected void updateCompanyBySymbol(Company company, CompanySymbols companySymbol, boolean shouldBeDeleted) {
|
|
||||||
assert company.getId().equals(companySymbol.getCompanyId());
|
|
||||||
String companySymbolValue = companySymbol.getCompanySymbolValue();
|
|
||||||
String setUpValue = shouldBeDeletedOrSet(companySymbolValue, shouldBeDeleted);
|
|
||||||
if (CompanySymbol.TRDC.equalsByKey(companySymbol.getCompanySymbol())) {
|
|
||||||
company.setTradingCode(setUpValue);
|
|
||||||
}
|
|
||||||
if (CompanySymbol.CLRC.equalsByKey(companySymbol.getCompanySymbol())) {
|
|
||||||
company.setClearingCode(setUpValue);
|
|
||||||
}
|
|
||||||
if (CompanySymbol.RGRC.equalsByKey(companySymbol.getCompanySymbol())) {
|
|
||||||
company.setRegistrationCode(setUpValue);
|
|
||||||
}
|
|
||||||
if (CompanySymbol.TAXN.equalsByKey(companySymbol.getCompanySymbol())) {
|
|
||||||
company.getProfile().setTaxNumber(setUpValue);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private String shouldBeDeletedOrSet(String companySymbol, boolean shouldBeDeleted) {
|
|
||||||
return shouldBeDeleted ? "" : companySymbol;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -0,0 +1,73 @@
|
||||||
|
package ru.spcex.clearing.fix.service;
|
||||||
|
|
||||||
|
import java.util.Collection;
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
|
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.company.CompanyRoleSet;
|
||||||
|
import ru.clearing.classes.statics.data.misc.OrderCurrency;
|
||||||
|
import ru.spcex.clearing.fix.quickfix.Application;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
|
import ru.spcex.platform.enumeration.CompanyRole;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||||
|
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class FixTaskListener extends QueueConsumer implements InitializingBean {
|
||||||
|
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final ImdgProvider imdgProvider;
|
||||||
|
private final Application fixApplication;
|
||||||
|
private final Imdg<CompanyRoleSet> companyRoleSetImdg;
|
||||||
|
private final Imdg<OrderCurrency> orderCurrencyImdg;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public FixTaskListener(Consumer<String, Object> kafkaQueue,
|
||||||
|
Producer<String, Object> kafkaProducer,
|
||||||
|
ImdgProvider imdgProvider,
|
||||||
|
Application fixApplication) {
|
||||||
|
super(kafkaQueue, kafkaProducer);
|
||||||
|
this.imdgProvider = imdgProvider;
|
||||||
|
this.fixApplication = fixApplication;
|
||||||
|
this.companyRoleSetImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanyRoleSet, CompanyRoleSet.class);
|
||||||
|
this.orderCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_OrderCurrency, OrderCurrency.class);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() {
|
||||||
|
callback(LauncherCommandRequest.class)
|
||||||
|
.setConsumer(this::processSendOrders)
|
||||||
|
.forDestination(Consts.SDOR_FIX_TASK, callbacks::put);
|
||||||
|
init();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void processSendOrders(BaseRequest<LauncherCommandRequest> request) {
|
||||||
|
log.info("SDOR[send order] started, request={}", request);
|
||||||
|
ImdgPredicateBuilder builder = orderCurrencyImdg.predicateBuilder();
|
||||||
|
CompanyRoleSet companyRoleSet = companyRoleSetImdg.getFirstObjectByPredicate(
|
||||||
|
builder.equals("companyRole", CompanyRole.CLCC.getKey())
|
||||||
|
);
|
||||||
|
|
||||||
|
ImdgPredicate clccOrderPredicate = builder.and(
|
||||||
|
builder.equals("status", "CRET"),
|
||||||
|
builder.or(
|
||||||
|
builder.equals("companyId", companyRoleSet.getCompanyId()),
|
||||||
|
builder.equals("counterPartyId", companyRoleSet.getCompanyId())
|
||||||
|
)
|
||||||
|
);
|
||||||
|
Collection<OrderCurrency> orderCurrencies = orderCurrencyImdg.getCollectionObjectsByPredicate(clccOrderPredicate);
|
||||||
|
orderCurrencies.forEach(fixApplication::sendOrder);
|
||||||
|
log.info("SDOR[send order] completed, count orders:{}", orderCurrencies.size());
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue