diff --git a/clearing-parent/default-management/src/main/java/ru/spcex/clearing/dflt/management/listener/TaskListener.java b/clearing-parent/default-management/src/main/java/ru/spcex/clearing/dflt/management/listener/TaskListener.java index c289ed898..efe2792a0 100644 --- a/clearing-parent/default-management/src/main/java/ru/spcex/clearing/dflt/management/listener/TaskListener.java +++ b/clearing-parent/default-management/src/main/java/ru/spcex/clearing/dflt/management/listener/TaskListener.java @@ -114,7 +114,7 @@ public class TaskListener extends QueueConsumer implements InitializingBean { } private void exportToTriAndSendToFix(BaseRequest req) { - + log.info("SDOR export and send command to fix, request={}", req); try { Collection orders = imdgQueryService.getClccOrderCurrencyIds(); String content = orderCurrencyTriExportService.exportToTri(orders); diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixService.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixService.java deleted file mode 100644 index 35437fb84..000000000 --- a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixService.java +++ /dev/null @@ -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 companyIMap; - protected UserRoleVerification userRoleVerification; -// protected CompanySymbolService companySymbolService; -// protected AccountNotificationHelper accountNotification; -// protected RelationService relationService; -// private final ValidationHelper validationHelper; -// private final Function companyNewRequestValidator; -// private final Function companyUpdateRequestValidator; -// private final Function companyDeleteRequestValidator; - - @Autowired - public FixService(Consumer kafkaQueue, Producer kafkaProducer, - ImdgProvider imdgProvider, - RequestHelper requestHelper, -// ValidationHelper validationHelper, - UserRoleVerification userRoleVerification -// @Qualifier("companyNewRequestValidator") -// Function companyNewRequestValidator, -// @Qualifier("companyUpdateRequestValidator") -// Function companyUpdateRequestValidator, -// @Qualifier("CompanyDeleteRequestValidator") -// Function 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 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 error = validator.tillFirstError(); -// if (error.isPresent()) { -// throw new ValidationException(error.get()); -// } -// } - ImdgTransaction transaction = imdgProvider.newTransaction(); - transaction.beginTransaction(); - boolean txOk = false; - try { - Imdg 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; - } -} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixTaskListener.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixTaskListener.java new file mode 100644 index 000000000..5cd243834 --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixTaskListener.java @@ -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 companyRoleSetImdg; + private final Imdg orderCurrencyImdg; + + @Autowired + public FixTaskListener(Consumer kafkaQueue, + Producer 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 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 orderCurrencies = orderCurrencyImdg.getCollectionObjectsByPredicate(clccOrderPredicate); + orderCurrencies.forEach(fixApplication::sendOrder); + log.info("SDOR[send order] completed, count orders:{}", orderCurrencies.size()); + } + +}