diff --git a/clearing-parent/fix-service/pom.xml b/clearing-parent/fix-service/pom.xml new file mode 100644 index 000000000..a990906fe --- /dev/null +++ b/clearing-parent/fix-service/pom.xml @@ -0,0 +1,131 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-3.22.285 + + 4.0.0 + + fix-service + + + 17 + 17 + + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.clearing + classes + + + ru.spcex.platform + platform-enum + + + ru.spcex.clearing + security-util + + + ru.spcex.clearing + clearing-validation + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + + + org.quickfixj + quickfixj-core + + + org.quickfixj + quickfixj-messages-fix44 + + + org.quickfixj + quickfixj-messages-fixt11 + + + + + net.logstash.logback + logstash-logback-encoder + 7.0.1 + + + org.springframework.boot + spring-boot-starter-test + test + + + ru.spcex.clearing + test-clearing + test + + + + + + + src/main/resources + + application.properties + fix.cfg + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + org.apache.maven.plugins + maven-surefire-plugin + 2.21.0 + + + org.junit.platform + junit-platform-surefire-provider + 1.2.0-M1 + + + org.junit.jupiter + junit-jupiter-engine + 5.2.0-M1 + + + + + + + + \ No newline at end of file diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/FixServiceApplication.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/FixServiceApplication.java new file mode 100644 index 000000000..7d639e578 --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/FixServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.fix; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class FixServiceApplication { + public static void main(String[] args) { + SpringApplication app = new SpringApplication(FixServiceApplication.class); + app.run(args); + } +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/BeanConfiguration.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/BeanConfiguration.java new file mode 100644 index 000000000..7c79f6aff --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/BeanConfiguration.java @@ -0,0 +1,33 @@ +package ru.spcex.clearing.fix.config; + +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; +import ru.spcex.clearing.fix.error.FixErrors; +import ru.spcex.clearing.util.security.UserRoleVerification; +import ru.spcex.clearing.util.services.IMDGMessageResolver; +import ru.spcex.clearing.util.services.RequestHelper; +import ru.spcex.platform.enumeration.UserRole; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +@Configuration +public class BeanConfiguration { + @Bean + public IMessageResolver messageResolver(ImdgProvider imdgProvider) { + return new IMDGMessageResolver(imdgProvider); + } + + @Bean + @Scope(value = ConfigurableBeanFactory.SCOPE_PROTOTYPE) + public RequestHelper requestHelper(IMessageResolver messageResolver) { + return new RequestHelper(messageResolver, FixErrors.GeneralError); + } + + @Bean + @Scope(value = ConfigurableBeanFactory.SCOPE_PROTOTYPE) + public UserRoleVerification userRoleVerification(ImdgProvider imdgProvider, IMessageResolver messageResolver) { + return new UserRoleVerification(imdgProvider, messageResolver, UserRole.Admin, FixErrors.UserVerifyDenial); + } +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/FixServiceImdgConfig.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/FixServiceImdgConfig.java new file mode 100644 index 000000000..6054f62af --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/FixServiceImdgConfig.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.fix.config; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import ru.spcex.clearing.fix.config.settings.FixServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class FixServiceImdgConfig { + private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) { + ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor(); + if (maxPoolSz > 2) { + pool.setKeepAliveSeconds(60); + pool.setAllowCoreThreadTimeOut(true); + } + pool.setCorePoolSize(maxPoolSz); + pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion); + return pool; + } + + @Bean(name = "taskExecutorHazelcastClientInitializer") + public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() { + return createThreadPoolTaskExecutor(1, true); + } + + @Bean(name = "taskExecutorIdGeneratorAwaiter") + public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() { + return createThreadPoolTaskExecutor(1, false); + } + + @Autowired + @Bean + public ImdgProvider imdgProvider( + @Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer, + @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter, + FixServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/KafkaConfig.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/KafkaConfig.java new file mode 100644 index 000000000..2940d5e08 --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/KafkaConfig.java @@ -0,0 +1,65 @@ +package ru.spcex.clearing.fix.config; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Scope; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.fix.config.settings.FixServiceSettings; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; +import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Configuration +public class KafkaConfig { + @Autowired + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean("kafkaConsumer") + public Consumer createConsumer(FixServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean("kafkaProducer") + public Producer createProducer(FixServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + + @Bean + public ProducerFactory pf(FixServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate, + ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } + +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/QuickFixConfig.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/QuickFixConfig.java new file mode 100644 index 000000000..375f8fb48 --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/QuickFixConfig.java @@ -0,0 +1,35 @@ +package ru.spcex.clearing.fix.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import quickfix.*; +import quickfix.field.ApplVerID; +import ru.spcex.clearing.fix.config.settings.FixServiceSettings; + +@Configuration +public class QuickFixConfig { + @Bean + public DataDictionary dataDictionary() throws ConfigError { + return new DataDictionary("FIX44.xml"); + } + + @Bean + public SessionSettings sessionSettings(FixServiceSettings settings) throws ConfigError { + return new SessionSettings(settings.getFixConfig()); + } + + @Bean + public MessageStoreFactory messageStoreFactory(SessionSettings sessionSettings) { + return new FileStoreFactory(sessionSettings); + } + + @Bean + public LogFactory logFactory(SessionSettings sessionSettings) { + return new FileLogFactory(sessionSettings); + } + + @Bean + public MessageFactory messageFactory() { + return new DefaultMessageFactory(ApplVerID.FIX44); + } +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/settings/FixServiceSettings.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/settings/FixServiceSettings.java new file mode 100644 index 000000000..02f50f81c --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/config/settings/FixServiceSettings.java @@ -0,0 +1,51 @@ +package ru.spcex.clearing.fix.config.settings; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.PropertySource; +import org.springframework.stereotype.Component; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams; + +@Component +@PropertySource("file:${spring.config.location}/application.properties") +@ConfigurationProperties("fix-service") +public class FixServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; + + private String fixConfig; + + public HazelcastClientParams getHazelcast() { + return hazelcast; + } + + public void setHazelcast(HazelcastClientParams hazelcast) { + this.hazelcast = hazelcast; + } + + public KafkaConsumerSettings getKafkaConsumer() { + return kafkaConsumer; + } + + public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) { + this.kafkaConsumer = kafkaConsumer; + } + + public KafkaProducerSettings getKafkaProducer() { + return kafkaProducer; + } + + public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { + this.kafkaProducer = kafkaProducer; + } + + public String getFixConfig() { + return fixConfig; + } + + public void setFixConfig(String fixConfig) { + this.fixConfig = fixConfig; + } +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/error/FixErrors.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/error/FixErrors.java new file mode 100644 index 000000000..328361b76 --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/error/FixErrors.java @@ -0,0 +1,25 @@ +package ru.spcex.clearing.fix.error; + +import ru.spcex.platform.utils.enumeration.IErrorEnumId; + +public enum FixErrors implements IErrorEnumId { + GeneralError(6000L), // Общая ошибка модуля fix-service + UserVerifyDenial(6001L), // 'Нет прав на проведение данной операции + + RequiredFieldEmpty(6002L), // Не заданы обязательные поля %s. + DictionaryNotFound(6003L), // Запись с идентификатором %s в справочнике %s не найдена. + WrongFieldValue(6004L), // Неверное значение поля %s. + RecordAlreadyExist(6005L), // Запись с идентификатором %s уже существует. + RecordNotFound(6006L); // Запись с идентификатором %s не найдена. + + private final Long id; + + FixErrors(Long id) { + this.id = id; + } + + @Override + public Long getId() { + return id; + } +} diff --git a/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/quickfix/Application.java b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/quickfix/Application.java new file mode 100644 index 000000000..1e8585ddf --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/quickfix/Application.java @@ -0,0 +1,191 @@ +package ru.spcex.clearing.fix.quickfix; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import quickfix.*; +import quickfix.field.*; +import quickfix.fix44.ExecutionReport; +import quickfix.fix44.NewOrderSingle; +import ru.clearing.classes.statics.data.misc.OrderCurrency; +import ru.clearing.classes.statics.data.misc.SOrders; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.ImdgTransaction; + +import java.math.BigDecimal; +import java.time.ZoneOffset; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +@Service +public class Application extends MessageCracker implements quickfix.Application, InitializingBean, DisposableBean { + private static final Logger log = LoggerFactory.getLogger(Application.class); + + private final ImdgProvider imdgProvider; + + private final SocketInitiator initiator; + + private SessionID fixSessionID; + private SessionID fdcSessionID; + + @Autowired + public Application(ImdgProvider imdgProvider, MessageStoreFactory storeFactory, SessionSettings settings, LogFactory logFactory, MessageFactory messageFactory) throws ConfigError { + this.imdgProvider = imdgProvider; + initiator = new SocketInitiator(this, storeFactory, settings, logFactory, messageFactory); + } + + @Override + public void afterPropertiesSet() throws ConfigError { + initiator.start(); + } + + @Override + public void destroy() throws Exception { + initiator.stop(); + } + + @Override + public void onCreate(SessionID sessionID) { + log.trace("onCreate {}", sessionID); + + if ("FIX".equals(sessionID.getSessionQualifier())) + fixSessionID = sessionID; + else if ("FDC".equals(sessionID.getSessionQualifier())) + fdcSessionID = sessionID; + } + + @Override + public void onLogon(SessionID sessionID) { + log.trace("onLogon {}", sessionID); + } + + @Override + public void onLogout(SessionID sessionID) { + log.trace("onLogout {}", sessionID); + } + + @Override + public void toAdmin(Message message, SessionID sessionID) { + log.trace("toAdmin {}, [{}]: {}", sessionID, message.getClass(), message); +// if (message instanceof Logon) { +// Logon msg = (Logon) message; +// msg.set(new Username(login)); +// msg.set(new Password(password)); +// msg.set(new EncryptMethod(EncryptMethod.NONE_OTHER)); +// msg.set(new HeartBtInt(30)); +// msg.set(new ResetSeqNumFlag(true)); +// } + } + + @Override + public void fromAdmin(Message message, SessionID sessionID) throws FieldNotFound, IncorrectDataFormat, IncorrectTagValue, RejectLogon { + log.trace("fromApp {}, [{}]: {}", sessionID, message.getClass(), message); + } + + @Override + public void toApp(Message message, SessionID sessionID) throws DoNotSend { + log.trace("toApp {}, [{}]: {}", sessionID, message.getClass(), message); + } + + @Override + public void fromApp(Message message, SessionID sessionID) throws FieldNotFound, IncorrectDataFormat, IncorrectTagValue, UnsupportedMessageType { + log.trace("fromApp {}, [{}]: {}", sessionID, message.getClass(), message); + } + + public void onMessage(ExecutionReport executionReport) { + ImdgTransaction transaction = imdgProvider.newTransaction(); + transaction.beginTransaction(); + try { + if (executionReport.isSetField(OrderID.FIELD) && executionReport.isSetField(ExDestination.FIELD) && executionReport.isSetField(TransactTime.FIELD)) { + OrderID orderId = executionReport.getOrderID(); + ExDestination exDestination = (ExDestination) executionReport.getField(new ExDestination()); + TransactTime transactTime = executionReport.getTransactTime(); + SecurityID securityID = executionReport.getSecurityID(); + Price price = executionReport.getPrice(); + OrderQty orderQty = executionReport.getOrderQty(); + Account account = executionReport.getAccount(); + Side side = executionReport.getSide(); + Currency currency = executionReport.getCurrency(); + SettlType settlType = executionReport.getSettlType(); + OrdStatus ordStatus = executionReport.getOrdStatus(); + + Pattern orderIdPattern = Pattern.compile("(\\d{6})-([A-Z]+)-(\\d+)"); + Matcher matcher = orderIdPattern.matcher(orderId.getValue()); + if (matcher.matches()) { +// LocalDate tradeDate = LocalDate.parse(matcher.group(1), DateTimeFormatter.ofPattern("yyMMdd")); +// String classCode = matcher.group(2); + Long orderNum = Long.parseLong(matcher.group(3)); + + Imdg sOrdersImdg = transaction.getImdg(IMDGDistributedNames.Map_SOrders, SOrders.class); + SOrders sOrders = sOrdersImdg.getSingleObjectByPredicate(sOrdersImdg.predicateBuilder().and( + sOrdersImdg.predicateBuilder().equals("orderNum", orderNum), + sOrdersImdg.predicateBuilder().equals("classCode", exDestination.getValue()), + sOrdersImdg.predicateBuilder().equals("tradeDate", transactTime.getValue().toLocalDate()))); + if (sOrders == null) { + sOrders = new SOrders(); + sOrders.setId(imdgProvider.getImdgIdGenerator().nextId()); + sOrders.setOrderNum(orderNum); + sOrders.setClassCode(exDestination.getValue()); + sOrders.setTradeDate(transactTime.getValue().toLocalDate()); + + sOrders.setSecCode(securityID.getValue()); + sOrders.setPrice(BigDecimal.valueOf(price.getValue())); + sOrders.setQty(BigDecimal.valueOf(orderQty.getValue())); + sOrders.setAccount(account.getValue()); + sOrders.setSettleCode(settlType.getValue()); + sOrders.setPriceCurrency(currency.getValue()); + + if (ordStatus.valueEquals(OrdStatus.NEW)) { + sOrders.setOrderDateTime(transactTime.getValue().toInstant(ZoneOffset.UTC)); + } else if (ordStatus.valueEquals(OrdStatus.CANCELED)) { + sOrders.setWithdrawDateTime(transactTime.getValue().toInstant(ZoneOffset.UTC)); + } + + sOrdersImdg.insert(sOrders); + } else { + if (ordStatus.valueEquals(OrdStatus.NEW)) { + sOrders.setOrderDateTime(transactTime.getValue().toInstant(ZoneOffset.UTC)); + } else if (ordStatus.valueEquals(OrdStatus.CANCELED)) { + sOrders.setWithdrawDateTime(transactTime.getValue().toInstant(ZoneOffset.UTC)); + } + + sOrdersImdg.update(sOrders); + } + } + } + transaction.commitTransaction(); + } catch (Exception e) { + log.error(executionReport.toString(), e); + transaction.rollbackTransaction(); + } + } + + public void sendOrder(OrderCurrency orderCurrency) { + try { + NewOrderSingle newOrderSingle = new NewOrderSingle(); + newOrderSingle.set(new ClOrdID(String.valueOf(orderCurrency.getId()))); + newOrderSingle.set(new Side(Side.BUY)); + newOrderSingle.set(new OrdType('R')); + newOrderSingle.set(new Price(orderCurrency.getPrice().doubleValue())); + newOrderSingle.set(new Account(orderCurrency.getAccount())); + newOrderSingle.set(new SecurityIDSource(SecurityIDSource.EXCHANGE_SYMBOL)); + newOrderSingle.set(new SecurityID(orderCurrency.getSecuritySymbol())); + newOrderSingle.set(new ExDestination(orderCurrency.getMarket())); + newOrderSingle.set(new OrderQty(orderCurrency.getQuantity().doubleValue())); + newOrderSingle.setField(new StringField(109, orderCurrency.getAccount())); // ClientID + newOrderSingle.set(new Text(orderCurrency.getComment())); + newOrderSingle.setField(new ContraBroker(orderCurrency.getCounterTradingCode())); + newOrderSingle.set(new SettlType(orderCurrency.getSettleCode())); + newOrderSingle.setField(new StringField(5018, orderCurrency.getTradingCode())); // FirmId + + boolean result = Session.sendToTarget(newOrderSingle, fixSessionID); + } catch (SessionNotFound sessionNotFound) { + log.error(orderCurrency.toString(), sessionNotFound); + } + } +} 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 new file mode 100644 index 000000000..35437fb84 --- /dev/null +++ b/clearing-parent/fix-service/src/main/java/ru/spcex/clearing/fix/service/FixService.java @@ -0,0 +1,197 @@ +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/resources/application.properties b/clearing-parent/fix-service/src/main/resources/application.properties new file mode 100644 index 000000000..e052708cf --- /dev/null +++ b/clearing-parent/fix-service/src/main/resources/application.properties @@ -0,0 +1,20 @@ +spring.main.web-application-type=none +fix-service.hazelcast.cluster-members=10.200.200.181:5701 +fix-service.hazelcast.login=dev +fix-service.hazelcast.password=dev-pass +fix-service.kafka-consumer.bootstrap-servers=10.200.200.181:9092 +fix-service.kafka-consumer.group-id=dev-group-fix-service +fix-service.kafka-consumer.enable-auto-commit=true +fix-service.kafka-consumer.session-timeout-ms=30000 +fix-service.kafka-consumer.auto-offset-reset=latest +fix-service.kafka-consumer.linger-ms=1 +fix-service.kafka-consumer.buffer-memory=33554432 + +fix-service.kafka-producer.bootstrap-servers=10.200.200.181:9092 +fix-service.kafka-producer.acks=all +fix-service.kafka-producer.retries=0 +fix-service.kafka-producer.batch-size=16384 +fix-service.kafka-producer.linger-ms=1 +fix-service.kafka-producer.buffer-memory=33554432 + +fix-service.fix-config=fix.cfg diff --git a/clearing-parent/fix-service/src/main/resources/fix.cfg b/clearing-parent/fix-service/src/main/resources/fix.cfg new file mode 100644 index 000000000..fd4120c63 --- /dev/null +++ b/clearing-parent/fix-service/src/main/resources/fix.cfg @@ -0,0 +1,27 @@ +[DEFAULT] +ConnectionType=initiator +SocketConnectHost=192.168.90.64 +StartTime=00:00:00 +EndTime=00:00:00 +HeartBtInt=30 +ReconnectInterval=5 +MaxLatency=1000 +ResetOnLogon=Y +ResetOnLogout=Y +FileStorePath=temp/fix +FileLogPath=temp/fix +ValidateFieldsHaveValues=N + +[SESSION] +BeginString=FIX.4.4 +SocketConnectPort=20038 +SenderCompID=KC_S968 +TargetCompID=SPCEX_UAT +SessionQualifier=FIX + +[SESSION] +BeginString=FIX.4.4 +SocketConnectPort=30020 +SenderCompID=KC_P47 +TargetCompID=SPCEX_UAT +SessionQualifier=FDC \ No newline at end of file diff --git a/clearing-parent/fix-service/src/main/resources/logback.xml b/clearing-parent/fix-service/src/main/resources/logback.xml new file mode 100644 index 000000000..3f4892ffc --- /dev/null +++ b/clearing-parent/fix-service/src/main/resources/logback.xml @@ -0,0 +1,56 @@ + + + + + + + + + + ${CONSOLE_LOG_PATTERN} + utf-8 + + + + + ${LOG_PATH}/${FILE_NAME}-text.log + + + ${FILE_LOG_PATTERN} + utf8 + + + ${LOG_PATH}/${FILE_NAME}-text.%d{yyyy-MM-dd}.%i.gz + + + 100MB + + 10 + + + + + ${LOG_PATH}/${FILE_NAME}-json.log + + + ${LOG_PATH}/${FILE_NAME}-json.%d{yyyy-MM-dd}.%i.gz + + + 100MB + + 10 + + + + + + + + + + + + + + + diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 37b018e3e..dec5eb5c1 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -34,6 +34,7 @@ scheduler-service clearing-service registry-service + fix-service test-clearing cleaning-builders trade-importer diff --git a/pom.xml b/pom.xml index ed0200dd7..9a8a4fe97 100644 --- a/pom.xml +++ b/pom.xml @@ -52,6 +52,7 @@ ${folder_root_clearing}/clearing-parent/utility-service ${folder_root_clearing}/clearing-parent/registry-service ${folder_root_clearing}/clearing-parent/scheduler-service + ${folder_root_clearing}/clearing-parent/fix-service ${folder_root_clearing}/clearing-parent/gateway-api ${folder_root_clearing}/clearing-parent/xml-exporter ${folder_root_clearing}/clearing-parent/xml-importer @@ -62,6 +63,7 @@ 2.8.1 2.9.2 1.5.22 + 2.3.1 @@ -197,6 +199,21 @@ swagger-annotations ${external_libraries.swagger.annotations.version} + + org.quickfixj + quickfixj-core + ${external_libraries.quickfix.version} + + + org.quickfixj + quickfixj-messages-fix44 + ${external_libraries.quickfix.version} + + + org.quickfixj + quickfixj-messages-fixt11 + ${external_libraries.quickfix.version} +