diff --git a/clearing-parent/backend-api/src/main/resources/meta/meta.xml b/clearing-parent/backend-api/src/main/resources/meta/meta.xml index 234da06e6..877c67b00 100644 --- a/clearing-parent/backend-api/src/main/resources/meta/meta.xml +++ b/clearing-parent/backend-api/src/main/resources/meta/meta.xml @@ -1,6 +1,6 @@ - + @@ -29,32 +29,32 @@ - + - + - + - - + + - + - - + + @@ -240,6 +240,11 @@ + + + + + @@ -288,7 +293,7 @@ - + @@ -302,41 +307,41 @@ - + - + - - + + - + - + - + - + - + - - + + @@ -347,7 +352,7 @@ - + @@ -355,8 +360,8 @@ - - + + @@ -364,12 +369,12 @@ - - + + - - + + @@ -377,10 +382,10 @@ - - - - + + + + @@ -396,8 +401,8 @@ - - + + @@ -408,12 +413,12 @@ - + - + @@ -506,7 +511,7 @@ - + @@ -852,107 +857,123 @@ - - - + + + - + + - + + + + + + + - - - + + + + + + + + + + + + + + + + + + + + + - - - - - - - - + + + + + + - - - - - - - - - - - - - - - - - - + - - - - - + + + + + - - - - - - - - + + + + + + + + + + + + + - - - - - + + + + + - - - - - - - - + + + + + + + + + + + + + - - - - - + + + + + - - - - - - - - - + + + + + + + + + @@ -1002,8 +1023,9 @@ + - + @@ -1021,8 +1043,9 @@ + - + @@ -1102,14 +1125,14 @@ - + - + @@ -1272,8 +1295,9 @@ - - + + + @@ -1324,7 +1348,7 @@ - + @@ -1333,7 +1357,7 @@
- + @@ -1352,7 +1376,7 @@ - + @@ -1415,8 +1439,9 @@ - + + @@ -1474,7 +1499,7 @@ - + @@ -1488,7 +1513,7 @@ - + @@ -1613,16 +1638,16 @@ - + - + - + @@ -1631,7 +1656,7 @@ - + @@ -1640,7 +1665,7 @@ - + @@ -1666,7 +1691,7 @@ - + @@ -1709,6 +1734,6 @@ - + diff --git a/clearing-parent/clearing-service/pom.xml b/clearing-parent/clearing-service/pom.xml new file mode 100644 index 000000000..aa206dcf4 --- /dev/null +++ b/clearing-parent/clearing-service/pom.xml @@ -0,0 +1,82 @@ + + + + clearing-parent + ru.spcex.clearing + SPCEX-1.0.0.0 + + 4.0.0 + + clearing-service + + + 17 + 17 + + + + + ru.spcex.platform + platform-messaging + + + ru.spcex.platform + platform-imdg-api-hazelcast-impl + + + ru.spcex.clearing + classes + + + org.springframework.boot + spring-boot-starter + + + com.fasterxml.jackson.core + jackson-databind + + + ru.spcex.platform + platform-enum + + + org.junit.jupiter + junit-jupiter + test + + + org.assertj + assertj-core + test + + + + + + src/main/resources + + application.properties + + false + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + repackage + + + + + ${project.artifactId} + + + + + \ No newline at end of file diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java new file mode 100644 index 000000000..3a4c724f8 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/ClearingServiceApplication.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class ClearingServiceApplication { + public static void main(String[] args) { + SpringApplication app = new SpringApplication(ClearingServiceApplication.class); + app.run(args); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java new file mode 100644 index 000000000..670a536a7 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ClearingImdgConfig.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.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.config.element.ClearingServiceSettings; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; + +@Configuration +public class ClearingImdgConfig { + 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, + ClearingServiceSettings settings + ) { + return new HazelcastService(taskExecutorHazelcastClientInitializer, + taskExecutorIdGeneratorAwaiter, + settings.getHazelcast()); + } + +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java new file mode 100644 index 000000000..0a281b5e1 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java @@ -0,0 +1,49 @@ +package ru.spcex.clearing.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 ru.spcex.clearing.config.element.ClearingServiceSettings; +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.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 + public Consumer createConsumer(ClearingServiceSettings settings) { + return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + } + + @Autowired + @Bean + public Producer createProducer(ClearingServiceSettings settings) { + return KafkaProducerFactory.producer(settings.getKafkaProducer()); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(Producer kafkaProducer, ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .producer(kafkaProducer) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java new file mode 100644 index 000000000..09157e14d --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessageResolverConfig.java @@ -0,0 +1,18 @@ +package ru.spcex.clearing.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +import java.util.Arrays; + +@Configuration +public class MessageResolverConfig { + @Bean + public IMessageResolver messageResolver() { + return errorMessage -> { + if (errorMessage == null) return "null"; + return String.format("(%d) args %s", errorMessage.getSubject().getId(), Arrays.toString(errorMessage.getArgs())); + }; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java new file mode 100644 index 000000000..38d52587e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/MessagesConfig.java @@ -0,0 +1,29 @@ +package ru.spcex.clearing.config; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.support.ResourceBundleMessageSource; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.enumeration.SpringPropertiesMessageResolver; + +import java.util.Locale; + +//@Configuration +public class MessagesConfig { + @Bean("validation-error-messages") + public ResourceBundleMessageSource messages() { + ResourceBundleMessageSource source = new ResourceBundleMessageSource(); + source.setBasenames("messages/error"); + source.setUseCodeAsDefaultMessage(true); + source.setDefaultEncoding("utf8"); + source.setDefaultLocale(Locale.ROOT); + return source; + } + + @Bean + public IMessageResolver errorResolver(@Qualifier("validation-error-messages") ResourceBundleMessageSource messageBundle) { + SpringPropertiesMessageResolver resolver = new SpringPropertiesMessageResolver(messageBundle); + resolver.setLocale("ru"); + return resolver; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java new file mode 100644 index 000000000..d42e5fcc3 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java @@ -0,0 +1,41 @@ +package ru.spcex.clearing.config.element; + +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("clearing-service") +public class ClearingServiceSettings { + private HazelcastClientParams hazelcast; + private KafkaConsumerSettings kafkaConsumer; + private KafkaProducerSettings kafkaProducer; + + 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; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java new file mode 100644 index 000000000..24106521d --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java @@ -0,0 +1,19 @@ +package ru.spcex.clearing.error; + +import ru.spcex.platform.utils.enumeration.IEnumId; + +public enum ClearingError implements IEnumId { + CompanyCreditCheck(10012L), + CompanyDebitCheck(10013L), + ; + private final Long id; + + ClearingError(Long id) { + this.id = id; + } + + @Override + public Long getId() { + return id; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java new file mode 100644 index 000000000..1103e95a6 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ClearingService.java @@ -0,0 +1,123 @@ +package ru.spcex.clearing.service; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.sdf.SDf03; +import ru.clearing.classes.statics.data.sdf.SDf11; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.TransactionStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.AbstractMap; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.stream.Collectors; + +@Service +@EnableScheduling +public class ClearingService { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg paymentImdgs; + private final ExecutorService executor; + private final ImdgId idGenerator; + private final PaymentInstructionSorter senderGroupSorter; + private final Imdg sdf03Imdg; + private final Imdg sdf11Imdg; + private final KafkaSender kafkaSender; + + @Autowired + public ClearingService(ImdgProvider imdgProvider, PaymentInstructionSorter senderGroupSorter, KafkaSender kafkaSender) { + this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class); + this.sdf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class); + this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class); + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.senderGroupSorter = senderGroupSorter; + this.kafkaSender = kafkaSender; + this.executor = Executors.newSingleThreadExecutor(); + } + + @Scheduled(cron = "${clearing-service.scheduler.check-payment-instruction}") + public void run() { + executor.execute(this::createSdfFromPaymentInstructionSTLD); + } + + private void createSdfFromPaymentInstructionSTLD() { + final boolean[] anyError = {false}; + //generationId для созадаваемых Sdf03/Sdf11 + Long generationId = idGenerator.nextId(); + //выгружаем PaymentInstructions с нужным статусом + Map paymentBySender = paymentImdgs.getCollectionObjectsByFieldValues( + Map.of("transactionStatus", TransactionStatus.stld.getKey())) + .stream() + //группируем по компаниям (fixme sorted убрать?) + .sorted(Comparator.comparing(PaymentInstruction::getSenderId)) + .collect(Collectors.groupingBy(PaymentInstruction::getSenderId)) + .entrySet() + .stream() + //результатом работы senderGroupSorter будет Map PaymentBatchInfo> + //PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment + //тип ClearingMemberCategory + .map(entry -> { + Long senderId = entry.getKey(); + List pmtInstrcs = entry.getValue(); + PaymentBatchInfo senderInfo = senderGroupSorter.sortCompanyPayments(generationId, senderId, pmtInstrcs); + if (senderInfo.getError() != null) { + anyError[0] = true; + } + return new AbstractMap.SimpleEntry<>(senderId, senderInfo); + }) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + //save SDF03/SDF11 + for (var entry : paymentBySender.entrySet()) { + PaymentBatchInfo senderPayments = entry.getValue(); + saveSdfAnSendToKafka(senderPayments, generationId); + } + //update PaymentInstruction.transactionStatus + for (var entry : paymentBySender.entrySet()) { + PaymentBatchInfo batch = entry.getValue(); + //все PaymentInstruction.transactionStatus в batch с error != null + //уже проапдейтились в методе sortCompanyPayments + if (batch.getError() == null) { + TransactionStatus stat = anyError[0] ? TransactionStatus.notSent : TransactionStatus.sent; + batch.getOrderedPaymentInstructions() + .forEach(paymentInstruction -> { + paymentInstruction.setTransactionStatus(stat.getKey()); + paymentImdgs.update(paymentInstruction); + }); + } + } + + } + + private void saveSdfAnSendToKafka(PaymentBatchInfo batch, Long generationId) { + SdfClearingRequest kafkaMessage = new SdfClearingRequest(); + kafkaMessage.setGroupId(generationId); + switch (batch.getCategoryD()) { + case I -> { + batch.getOrderedPaymentInstructions() + .map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId)) + .forEach(sdf03Imdg::insert); + kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage); + } + case B -> { + batch.getOrderedPaymentInstructions() + .map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId)) + .forEach(sdf11Imdg::insert); + kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage); + } + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java new file mode 100644 index 000000000..fff82f55c --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java @@ -0,0 +1,69 @@ +package ru.spcex.clearing.service; + +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.spcex.platform.enumeration.ClearingMemberCategoryD; +import ru.spcex.platform.utils.enumeration.EnumMessage; + +import java.util.Collection; +import java.util.List; +import java.util.function.Function; +import java.util.stream.Stream; + +public class PaymentBatchInfo { + private List initialOrder; + private List fromClearingToBank; + private List fromBankToClearing; + private ClearingMemberCategoryD categoryD; + private EnumMessage error; + + public Stream getOrderedPaymentInstructions() { + Function, Stream> safeStream + = paymentInstructions -> paymentInstructions != null ? paymentInstructions.stream() : Stream.empty(); + if (categoryD.equals(ClearingMemberCategoryD.B)) { + return safeStream.apply(initialOrder); + } else { + return Stream.concat(safeStream.apply(fromClearingToBank), safeStream.apply(fromBankToClearing)); + } + } + + + public List getFromClearingToBank() { + return fromClearingToBank; + } + + public void setFromClearingToBank(List fromClearingToBank) { + this.fromClearingToBank = fromClearingToBank; + } + + public List getFromBankToClearing() { + return fromBankToClearing; + } + + public void setFromBankToClearing(List fromBankToClearing) { + this.fromBankToClearing = fromBankToClearing; + } + + public ClearingMemberCategoryD getCategoryD() { + return categoryD; + } + + public void setCategoryD(ClearingMemberCategoryD categoryD) { + this.categoryD = categoryD; + } + + public EnumMessage getError() { + return error; + } + + public void setError(EnumMessage error) { + this.error = error; + } + + public List getInitialOrder() { + return initialOrder; + } + + public void setInitialOrder(List initialOrder) { + this.initialOrder = initialOrder; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentInstructionSorter.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentInstructionSorter.java new file mode 100644 index 000000000..f5f59d15e --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentInstructionSorter.java @@ -0,0 +1,121 @@ +package ru.spcex.clearing.service; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.account.Account; +import ru.clearing.classes.statics.data.generated.ClearingMemberCategory; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.spcex.clearing.error.ClearingError; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.enumeration.AccountType; +import ru.spcex.platform.enumeration.ClearingMemberCategoryD; +import ru.spcex.platform.enumeration.TransactionStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumKey; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.function.Function; + +@Component +public class PaymentInstructionSorter { + + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg clrngMmbrImdg; + private final Imdg accImdg; + private final IMessageResolver messageResolver; + private final Imdg pmtInstrctnsImdg; + + @Autowired + public PaymentInstructionSorter(ImdgProvider imdgProvider, IMessageResolver messageResolver) { + this.clrngMmbrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); + this.pmtInstrctnsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class); + this.accImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + this.messageResolver = messageResolver; + } + + PaymentBatchInfo sortCompanyPayments(Long generationId, Long senderId, List payments) { + PaymentBatchInfo batchInfo = new PaymentBatchInfo(); + log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size()); + ClearingMemberCategory category = clrngMmbrImdg.getSingleObjectByFieldValues(Map.of("companyId", senderId)); + ClearingMemberCategoryD categoryValue = IEnumKey.getEnumByKey(ClearingMemberCategoryD.class, category.getClearingMemberCategory()); + batchInfo.setCategoryD(categoryValue); + if (ClearingMemberCategoryD.I.equals(categoryValue)) { + EnumMessage error = null; + List aList = new ArrayList<>(); + List bList = new ArrayList<>(); + for (PaymentInstruction payment : payments) { + Account creditLegAcc = accImdg.getSingleObjectByID(payment.getCreditLegAccountId()); + Account debitLegAcc = accImdg.getSingleObjectByID(payment.getDebitLegAccountId()); + if (creditLegAcc == null || debitLegAcc == null) { + String message = String.format("cannot find account CreditLegAccountId/DebitLegAccountId %d/%d", payment.getCreditLegAccountId(), payment.getDebitLegAccountId()); + log.error("generationId={}, senderId={} payment.id={} {}", + generationId, senderId, payment.getId(), message); + continue; + } + if (AccountType.Clrn.equalsByKey(creditLegAcc.getAccountType()) + && AccountType.Bank.equalsByKey(debitLegAcc.getAccountType())) { + aList.add(payment); + } else if (AccountType.Bank.equalsByKey(creditLegAcc.getAccountType()) + && AccountType.Clrn.equalsByKey(debitLegAcc.getAccountType())) { + bList.add(payment); + } else { + String message = String.format("cannot sort creditLegAccount.type=%s, debitLegAccount.type=%s", creditLegAcc.getAccountType(), debitLegAcc.getAccountType()); + log.error("generationId={}, senderId={} payment.id={} {}", + generationId, senderId, payment.getId(), message); + //continue; + } + } + Function, Long> creditAmountSum = paymentInstructions -> paymentInstructions + .stream() + .map(PaymentInstruction::getCreditLegAmount) + .reduce(0L, Long::sum); + Function, Long> debitAmountSum = paymentInstructions -> paymentInstructions + .stream() + .map(PaymentInstruction::getDebitLegAmount) + .reduce(0L, Long::sum); + Long fromClearingToBankCreditAmount = creditAmountSum.apply(aList); + Long fromBankToClearingCreditAmount = creditAmountSum.apply(bList); + Long fromClearingToBankDebitAmount = debitAmountSum.apply(aList); + Long fromBankToClearingDebitAmount = debitAmountSum.apply(bList); + if (!fromClearingToBankCreditAmount.equals(fromBankToClearingCreditAmount)) { + error = new EnumMessage(ClearingError.CompanyCreditCheck, senderId.toString()); + } else if (!fromClearingToBankDebitAmount.equals(fromBankToClearingDebitAmount)) { + error = new EnumMessage(ClearingError.CompanyDebitCheck, senderId.toString()); + } + log.info("processing PaymentInstruction's generationId={} senderId={} [fromClearingToBankCreditAmount={}, " + + "fromBankToClearingCreditAmount={}, " + + "fromClearingToBankDebitAmount={}, " + + "fromBankToClearingDebitAmount={}] {}", generationId, senderId, + fromClearingToBankCreditAmount, + fromBankToClearingCreditAmount, + fromClearingToBankDebitAmount, + fromBankToClearingDebitAmount, + error != null ? ("error " + messageResolver.resolve(error)) : "ok"); + if (error != null) { + batchInfo.setError(error); + payments.forEach(pmt -> { + pmt.setTransactionStatus(TransactionStatus.cher.getKey()); + pmtInstrctnsImdg.update(pmt); + }); + return batchInfo; + } else { + batchInfo.setFromClearingToBank(aList); + batchInfo.setFromBankToClearing(bList); + return batchInfo; + } + } else if (ClearingMemberCategoryD.B.equals(categoryValue)) { + batchInfo.setInitialOrder(payments); + return batchInfo; + } else { + throw new IllegalStateException("unknown clearing member category for generationId=" + + generationId + " companyId=" + senderId + ""); + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf03Builder.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf03Builder.java new file mode 100644 index 000000000..fc32df0e2 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf03Builder.java @@ -0,0 +1,53 @@ +package ru.spcex.clearing.service; + +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.sdf.SDf03; +import ru.spcex.platform.enumeration.Sender; +import ru.spcex.platform.utils.time.TimeUtil; + +import java.time.Instant; +import java.time.LocalDate; + +public class Sdf03Builder { + public static SDf03 buildSdf03(PaymentInstruction paymentInstruction, Long generationId) { + SDf03 sDf03 = new SDf03(); + sDf03.setSeg_type("S"); + sDf03.setDoc_type("002"); + sDf03.setDocnm_ref(paymentInstruction.getId().toString()); + sDf03.setPriority("9"); + sDf03.setSbankcode(Sender.Prc.equalsById(paymentInstruction.getSenderId()) ? "@09999" : "@0"); + sDf03.setC_acc_deb(paymentInstruction.getDebitLegAccount()); + sDf03.setSbanknam1(paymentInstruction.getPayeeBankName()); + // =@00100, если ?; + // =@00000, если ?. + sDf03.setRbankcode(Sender.Prc.equalsById(paymentInstruction.getSenderId()) ? "@09999" : "@0"); + sDf03.setC_acc_cred(paymentInstruction.getCreditLegAccount()); + sDf03.setRbanknam1(paymentInstruction.getAddresseeBankName()); + sDf03.setPay_date(LocalDate.now().format(TimeUtil.PROPERTY_DATE_FORMATTER)); + sDf03.setPay_val("RUR"); + sDf03.setSum_deb(paymentInstruction.getDebitLegAmount() != null ? paymentInstruction.getDebitLegAmount().toString() : null); + //37 sp_code varchar(2) Код назначения платежа + String[] splitPaymentPurpose = SpecifUtil.splitPaymentPurpose(paymentInstruction.getPaymentPurpose()); + for (int i = 0; i < splitPaymentPurpose.length; i++) { + String specif = splitPaymentPurpose[i]; + if (i == 0) { + sDf03.setSpecif_1(specif); + } else if (i == 1) { + sDf03.setSpecif_2(specif); + } else if (i == 2) { + sDf03.setSpecif_3(specif); + } else if (i == 3) { + sDf03.setSpecif_4(specif); + } else if (i == 4) { + sDf03.setSpecif_5(specif); + } else if (i == 5) { + sDf03.setSpecif_6(specif); + } + } + sDf03.setGenerationTime(Instant.now()); + sDf03.setGenerationId(generationId); + //todo +// sDf03.setPaymentInstructionId(paymentInstruction.getId()); + return sDf03; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf11Builder.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf11Builder.java new file mode 100644 index 000000000..2d37e878b --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf11Builder.java @@ -0,0 +1,53 @@ +package ru.spcex.clearing.service; + +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.sdf.SDf11; +import ru.spcex.platform.enumeration.Sender; +import ru.spcex.platform.utils.time.TimeUtil; + +import java.time.Instant; +import java.time.LocalDate; + +public class Sdf11Builder { + public static SDf11 buildSdf11(PaymentInstruction paymentInstruction, Long generationId) { + SDf11 sDf11 = new SDf11(); + sDf11.setSeg_type("S"); + sDf11.setDoc_type("002"); + sDf11.setDocnm_ref(paymentInstruction.getId().toString()); + sDf11.setPriority("9"); + sDf11.setSbankcode(Sender.Prc.equalsById(paymentInstruction.getSenderId()) ? "@09999" : "@0"); + sDf11.setC_acc_deb(paymentInstruction.getDebitLegAccount()); + sDf11.setSbanknam1(paymentInstruction.getPayeeBankName()); + // =@00100, если ?; + // =@00000, если ?. + sDf11.setRbankcode(Sender.Prc.equalsById(paymentInstruction.getSenderId()) ? "@09999" : "@0"); + sDf11.setC_acc_cred(paymentInstruction.getCreditLegAccount()); + sDf11.setRbanknam1(paymentInstruction.getAddresseeBankName()); + sDf11.setPay_date(LocalDate.now().format(TimeUtil.PROPERTY_DATE_FORMATTER)); + sDf11.setPay_val("RUR"); + sDf11.setSum_deb(paymentInstruction.getDebitLegAmount() != null ? paymentInstruction.getDebitLegAmount().toString() : null); + //37 sp_code varchar(2) Код назначения платежа + String[] splitPaymentPurpose = SpecifUtil.splitPaymentPurpose(paymentInstruction.getPaymentPurpose()); + for (int i = 0; i < splitPaymentPurpose.length; i++) { + String specif = splitPaymentPurpose[i]; + if (i == 0) { + sDf11.setSpecif_1(specif); + } else if (i == 1) { + sDf11.setSpecif_2(specif); + } else if (i == 2) { + sDf11.setSpecif_3(specif); + } else if (i == 3) { + sDf11.setSpecif_4(specif); + } else if (i == 4) { + sDf11.setSpecif_5(specif); + } else if (i == 5) { + sDf11.setSpecif_6(specif); + } + } + sDf11.setGenerationTime(Instant.now()); + sDf11.setGenerationId(generationId); + //todo +// sDf03.setPaymentInstructionId(paymentInstruction.getId()); + return sDf11; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SpecifUtil.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SpecifUtil.java new file mode 100644 index 000000000..79ac236da --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SpecifUtil.java @@ -0,0 +1,21 @@ +package ru.spcex.clearing.service; + +public class SpecifUtil { + public static String[] splitPaymentPurpose(String paymentPurpose) { + if (paymentPurpose == null || paymentPurpose.length() < 1) return new String[0]; + int specifSize = divisionRoundUp(paymentPurpose.length(), 35); + specifSize = Math.min(specifSize, 6); + String[] res = new String[specifSize]; + for (int i = 0; i < res.length; i++) { + int startIndex = i * 35; + int endIndex = Math.min(35 * (i + 1), paymentPurpose.length()); + res[i] = paymentPurpose.substring(startIndex, endIndex); + } + return res; + } + + public static int divisionRoundUp(int a, int b) { + int ostatok = a % b; + return a / b + (ostatok > 0 ? 1 : 0); + } +} diff --git a/clearing-parent/clearing-service/src/main/resources/application.properties b/clearing-parent/clearing-service/src/main/resources/application.properties new file mode 100644 index 000000000..af4ebcf34 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/resources/application.properties @@ -0,0 +1,22 @@ +spring.main.web-application-type=none + +clearing-service.hazelcast.cluster-members=127.0.0.1:5701 +clearing-service.hazelcast.login=dev +clearing-service.hazelcast.password=dev-pass + +clearing-service.kafka-consumer.bootstrap-servers=localhost:9092 +clearing-service.kafka-consumer.group-id=dev-group-clearing-service +clearing-service.kafka-consumer.enable-auto-commit=false +clearing-service.kafka-consumer.session-timeout-ms=30000 +clearing-service.kafka-consumer.auto-offset-reset=latest +clearing-service.kafka-consumer.linger-ms=1 +clearing-service.kafka-consumer.buffer-memory=33554432 + +clearing-service.kafka-producer.bootstrap-servers=localhost:9092 +clearing-service.kafka-producer.acks=all +clearing-service.kafka-producer.retries=0 +clearing-service.kafka-producer.batch-size=16384 +clearing-service.kafka-producer.linger-ms=1 +clearing-service.kafka-producer.buffer-memory=33554432 + +clearing-service.scheduler.check-payment-instruction=*/5 * * * * * \ No newline at end of file diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/service/SpecifUtilTest.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/service/SpecifUtilTest.java new file mode 100644 index 000000000..a53305668 --- /dev/null +++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/service/SpecifUtilTest.java @@ -0,0 +1,76 @@ +package ru.spcex.clearing.service; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +class SpecifUtilTest { + @Test + public void testDivideRoundUp() { + Assertions.assertEquals(2, SpecifUtil.divisionRoundUp(36, 35)); + Assertions.assertEquals(2, SpecifUtil.divisionRoundUp(70, 35)); + Assertions.assertEquals(3, SpecifUtil.divisionRoundUp(71, 35)); + } + + @Test + public void testSplit() { + { + //40 symbols + String[] splitted = SpecifUtil.splitPaymentPurpose("0123456789012345678901234567890123456789"); + Assertions.assertEquals(2, splitted.length); + Assertions.assertEquals("01234567890123456789012345678901234", splitted[0]); + Assertions.assertEquals("56789", splitted[1]); + } + + { + //40 symbols + String[] splitted = SpecifUtil.splitPaymentPurpose("01234567890123456789012345678901234"); + Assertions.assertEquals(1, splitted.length); + Assertions.assertEquals("01234567890123456789012345678901234", splitted[0]); + } + + { + //10 symbols + String[] splitted = SpecifUtil.splitPaymentPurpose("0123456789"); + Assertions.assertEquals(1, splitted.length); + Assertions.assertEquals("0123456789", splitted[0]); + } + + { + //0 symbols + String[] splitted = SpecifUtil.splitPaymentPurpose(""); + Assertions.assertEquals(0, splitted.length); + } + + { + //80 symbols + String[] splitted = SpecifUtil.splitPaymentPurpose("0123456789012345678901234567890123456789" + + "0123456789012345678901234567890123456789"); + Assertions.assertEquals(3, splitted.length); + Assertions.assertEquals("01234567890123456789012345678901234", splitted[0]); + Assertions.assertEquals("56789012345678901234567890123456789", splitted[1]); + Assertions.assertEquals("0123456789", splitted[2]); + } + + { + //maximum symbols + String[] splitted = SpecifUtil.splitPaymentPurpose("01234567890123456789012345678912345" + + "01234567890123456789012345678912345" + + "01234567890123456789012345678912345" + + "01234567890123456789012345678912345" + + "01234567890123456789012345678912345" + + "01234567890123456789012345678912345" + + "outside of scope" + ); + Assertions.assertEquals(6, splitted.length); + for (String s : splitted) { + Assertions.assertEquals("01234567890123456789012345678912345", s); + } + } + + + +// ""; + } + + +} \ No newline at end of file diff --git a/clearing-parent/pom.xml b/clearing-parent/pom.xml index 08cc8557b..85fb2a0a6 100644 --- a/clearing-parent/pom.xml +++ b/clearing-parent/pom.xml @@ -30,6 +30,7 @@ account-service balance-service scheduler-service + clearing-service diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java index effaaf686..b07db3270 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/SchedulerService.java @@ -7,73 +7,75 @@ 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.misc.KeyRate; +import ru.clearing.classes.statics.data.scheduler.Scheduler; 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.common.CommonDeleteRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateNewRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.utilities.KeyRateUpdateRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.SchedulerNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.SchedulerUpdateRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; -import ru.spcex.platform.enumeration.Status; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @Service public class SchedulerService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); - private final Imdg keyRateMap; + private final Imdg schedulerMap; @Autowired public SchedulerService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { super(kafkaQueue, kafkaProducer); - this.keyRateMap = imdgProvider.getImdg(IMDGDistributedNames.Map_KeyRate, KeyRate.class); + this.schedulerMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Scheduler, Scheduler.class); } @Override public void afterPropertiesSet() { - callback(KeyRateNewRequest.class) - .setConsumer(this::newKeyRate) - .forDestination(Consts.DESTINATION_KEY_RATE_NEW, callbacks::put); - callback(KeyRateUpdateRequest.class) - .setConsumer(this::updateKeyRate) - .forDestination(Consts.DESTINATION_KEY_RATE_UPDATE, callbacks::put); + callback(SchedulerNewRequest.class) + .setConsumer(this::newScheduler) + .forDestination(Consts.DESTINATION_SCHEDULER_NEW, callbacks::put); + callback(SchedulerUpdateRequest.class) + .setConsumer(this::updateScheduler) + .forDestination(Consts.DESTINATION_SCHEDULER_UPDATE, callbacks::put); callback(CommonDeleteRequest.class) - .setConsumer(this::deleteKeyRate) - .forDestination(Consts.DESTINATION_KEY_RATE_DELETE, callbacks::put); + .setConsumer(this::deleteScheduler) + .forDestination(Consts.DESTINATION_SCHEDULER_DELETE, callbacks::put); init(); } - private void newKeyRate(BaseRequest userRequest) { - KeyRateNewRequest req = userRequest.getRequestPayload(); - log.debug("MoneyMarketSecurityNewRequest received"); - KeyRate keyRate = new KeyRate(); - keyRate.setEndDate(req.getEndDate()); - keyRate.setDocument(req.getDocument()); - keyRate.setRate(req.getKeyRate()); - keyRate.setStartDate(req.getStartDate()); - keyRate.setWorkflowStatus(Status.Active.getKey()); - keyRateMap.insert(keyRate); - log.debug("successfully processed, new id {}", keyRate.getId()); + private void newScheduler(BaseRequest userRequest) { + SchedulerNewRequest req = userRequest.getRequestPayload(); + log.debug("SchedulerNewRequest received"); + Scheduler scheduler = new Scheduler(); + scheduler.setTask(req.getTask()); + scheduler.setTaskTime(req.getTaskTime()); + scheduler.setClearingDate(req.getClearingDate()); + scheduler.setMarket(req.getMarket()); + scheduler.setTaskStatus(req.getTaskStatus()); + scheduler.setSecurityId(req.getSecurityId()); + schedulerMap.insert(scheduler); + log.debug("successfully processed, new id {}", scheduler.getId()); } - private void updateKeyRate(BaseRequest userRequest) { - KeyRateUpdateRequest req = userRequest.getRequestPayload(); + private void updateScheduler(BaseRequest userRequest) { + SchedulerUpdateRequest req = userRequest.getRequestPayload(); log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId()); - KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId()); - keyRate.setEndDate(req.getEndDate()); - keyRate.setDocument(req.getDocument()); - keyRate.setRate(req.getKeyRate()); - keyRate.setStartDate(req.getStartDate()); - keyRateMap.update(keyRate); + Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId()); + scheduler.setTask(req.getTask()); + scheduler.setTaskTime(req.getTaskTime()); + scheduler.setClearingDate(req.getClearingDate()); + scheduler.setMarket(req.getMarket()); + scheduler.setTaskStatus(req.getTaskStatus()); + scheduler.setSecurityId(req.getSecurityId()); + schedulerMap.update(scheduler); } - private void deleteKeyRate(BaseRequest userRequest) { + private void deleteScheduler(BaseRequest userRequest) { CommonDeleteRequest req = userRequest.getRequestPayload(); log.debug("CommonDeleteRequest received id = {}", req.getId()); - KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId()); - keyRateMap.delete(keyRate); + Scheduler scheduler = schedulerMap.getSingleObjectByID(req.getId()); + schedulerMap.delete(scheduler); } } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TimetableService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TimetableService.java new file mode 100644 index 000000000..bc77439c0 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TimetableService.java @@ -0,0 +1,75 @@ +package ru.spcex.clearing.scheduler.service; + +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.scheduler.Timetable; +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.common.CommonDeleteRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TimetableNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TimetableUpdateRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class TimetableService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg timetableMap; + + @Autowired + public TimetableService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); + this.timetableMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Timetable, Timetable.class); + } + + @Override + public void afterPropertiesSet() { + callback(TimetableNewRequest.class) + .setConsumer(this::newTimetable) + .forDestination(Consts.DESTINATION_TIMETABLE_NEW, callbacks::put); + callback(TimetableUpdateRequest.class) + .setConsumer(this::updateTimetable) + .forDestination(Consts.DESTINATION_TIMETABLE_UPDATE, callbacks::put); + callback(CommonDeleteRequest.class) + .setConsumer(this::deleteTimetable) + .forDestination(Consts.DESTINATION_TIMETABLE_DELETE, callbacks::put); + init(); + } + + private void newTimetable(BaseRequest userRequest) { + TimetableNewRequest req = userRequest.getRequestPayload(); + log.debug("TimetableNewRequest received"); + Timetable timetable = new Timetable(); + timetable.setTask(req.getTask()); + timetable.setTaskTime(req.getTaskTime()); + timetable.setTaskStatus(req.getTaskStatus()); + timetableMap.insert(timetable); + log.debug("successfully processed, new id {}", timetable.getId()); + } + + private void updateTimetable(BaseRequest userRequest) { + TimetableUpdateRequest req = userRequest.getRequestPayload(); + log.debug("TimetableNewRequest received"); + Timetable timetable = timetableMap.getSingleObjectByID(req.getId()); + timetable.setTask(req.getTask()); + timetable.setTaskTime(req.getTaskTime()); + timetable.setTaskStatus(req.getTaskStatus()); + timetableMap.update(timetable); + log.debug("successfully processed, new id {}", timetable.getId()); + } + + private void deleteTimetable(BaseRequest userRequest) { + CommonDeleteRequest req = userRequest.getRequestPayload(); + log.debug("CommonDeleteRequest received id = {}", req.getId()); + Timetable timetable = timetableMap.getSingleObjectByID(req.getId()); + timetableMap.delete(timetable); + } +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TradingCalendarService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TradingCalendarService.java new file mode 100644 index 000000000..baf9cbda6 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TradingCalendarService.java @@ -0,0 +1,75 @@ +package ru.spcex.clearing.scheduler.service; + +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.scheduler.TradingCalendar; +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.common.CommonDeleteRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TradingCalendarNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.TradingCalendarUpdateRequest; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Service +public class TradingCalendarService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg tradingCalendarMap; + + @Autowired + public TradingCalendarService(Consumer kafkaQueue, Producer kafkaProducer, + ImdgProvider imdgProvider) { + super(kafkaQueue, kafkaProducer); + this.tradingCalendarMap = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingCalendar, TradingCalendar.class); + } + + @Override + public void afterPropertiesSet() { + callback(TradingCalendarNewRequest.class) + .setConsumer(this::newTradingCalendar) + .forDestination(Consts.DESTINATION_TRADING_CALENDAR_NEW, callbacks::put); + callback(TradingCalendarUpdateRequest.class) + .setConsumer(this::updateTradingCalendar) + .forDestination(Consts.DESTINATION_TRADING_CALENDAR_UPDATE, callbacks::put); + callback(CommonDeleteRequest.class) + .setConsumer(this::deleteTradingCalendar) + .forDestination(Consts.DESTINATION_TRADING_CALENDAR_DELETE, callbacks::put); + init(); + } + + private void newTradingCalendar(BaseRequest userRequest) { + TradingCalendarNewRequest req = userRequest.getRequestPayload(); + log.debug("TradingCalendarNewRequest received"); + TradingCalendar tradingCalendar = new TradingCalendar(); + tradingCalendar.setClearingDate(req.getClearingDate()); + tradingCalendar.setCompanyId(req.getCompanyId()); + tradingCalendar.setTradingStatus(req.getTradingStatus()); + tradingCalendarMap.insert(tradingCalendar); + log.debug("successfully processed, new id {}", tradingCalendar.getId()); + } + + private void updateTradingCalendar(BaseRequest userRequest) { + TradingCalendarUpdateRequest req = userRequest.getRequestPayload(); + log.debug("TradingCalendarUpdateRequest received"); + TradingCalendar tradingCalendar = tradingCalendarMap.getSingleObjectByID(req.getId()); + tradingCalendar.setClearingDate(req.getClearingDate()); + tradingCalendar.setCompanyId(req.getCompanyId()); + tradingCalendar.setTradingStatus(req.getTradingStatus()); + tradingCalendarMap.update(tradingCalendar); + log.debug("successfully processed, new id {}", tradingCalendar.getId()); + } + + private void deleteTradingCalendar(BaseRequest userRequest) { + CommonDeleteRequest req = userRequest.getRequestPayload(); + log.debug("CommonDeleteRequest received id = {}", req.getId()); + TradingCalendar tradingCalendar = tradingCalendarMap.getSingleObjectByID(req.getId()); + tradingCalendarMap.delete(tradingCalendar); + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java index 71385a733..bf0f5f50b 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/AccountType.java @@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum AccountType implements IEnumKey { - Clrn("CLRN"); + Clrn("CLRN"), Bank("BANK"); private final String key; diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java new file mode 100644 index 000000000..bc94d111c --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ClearingMemberCategoryD.java @@ -0,0 +1,18 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum ClearingMemberCategoryD implements IEnumKey { + I("I"), B("B"); + + private final String key; + + ClearingMemberCategoryD(String key) { + this.key = key; + } + + @Override + public String getKey() { + return key; + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java new file mode 100644 index 000000000..2d97429b9 --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/TransactionStatus.java @@ -0,0 +1,18 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum TransactionStatus implements IEnumKey { + stld("STLD"), notSent("NSNT"), cher("CHER"), sent("SENT"); + + private final String key; + + TransactionStatus(String key) { + this.key = key; + } + + @Override + public String getKey() { + return key; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index df5ac80b5..2e4b2721b 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -9,6 +9,14 @@ public interface Consts { String DESTINATION_KEY_RATE_UPDATE = "key-rate-update"; String DESTINATION_KEY_RATE_DELETE = "key-rate-delete"; + String DESTINATION_TIMETABLE_NEW = "timetable-new"; + String DESTINATION_TIMETABLE_UPDATE = "timetable-update"; + String DESTINATION_TIMETABLE_DELETE = "timetable-delete"; + + String DESTINATION_TRADING_CALENDAR_NEW = "timetable-new"; + String DESTINATION_TRADING_CALENDAR_UPDATE = "timetable-update"; + String DESTINATION_TRADING_CALENDAR_DELETE = "timetable-delete"; + String DESTINATION_SCHEDULER_NEW = "scheduler-new"; String DESTINATION_SCHEDULER_UPDATE = "scheduler-update"; String DESTINATION_SCHEDULER_DELETE = "scheduler-delete"; @@ -35,6 +43,8 @@ public interface Consts { //todo String STATEMENT_PROCESS = "statement-process"; String SDF04_PROCESS = "sdf04-process"; + String SDF03_PROCESS = "sdf03-process"; + String SDF11_PROCESS = "sdf11-process"; String EXPORT_PROCESS = "export-process"; String ACCOUNT_NEW = "account-new"; String BALANCE_ACCOUNT_NEW = "balance-account-new"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/SdfClearingRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/SdfClearingRequest.java new file mode 100644 index 000000000..164e9c15d --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/SdfClearingRequest.java @@ -0,0 +1,16 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.clearing; + +import com.fasterxml.jackson.annotation.JsonProperty; + +public class SdfClearingRequest { + @JsonProperty + private Long groupId; + + public Long getGroupId() { + return groupId; + } + + public void setGroupId(Long groupId) { + this.groupId = groupId; + } +} \ No newline at end of file diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java index 405d98f59..9b72792d3 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerNewRequest.java @@ -4,23 +4,26 @@ import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.databind.annotation.JsonDeserialize; import com.fasterxml.jackson.databind.annotation.JsonSerialize; import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer; import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer; import java.time.LocalDate; import java.time.LocalTime; public class SchedulerNewRequest { - @JsonProperty - @JsonSerialize(using = LocalDateSerializer.class) - @JsonDeserialize(using = LocalDateDeserializer.class) - public LocalTime taskTime; @JsonProperty public String task; + @JsonSerialize(using = LocalTimeSerializer.class) + @JsonDeserialize(using = LocalTimeDeserializer.class) @JsonProperty + public LocalTime taskTime; + @JsonSerialize(using = LocalDateSerializer.class) @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty public LocalDate clearingDate; @JsonProperty @@ -32,13 +35,12 @@ public class SchedulerNewRequest { @JsonProperty public Long securityId; - public SchedulerNewRequest(LocalTime taskTime, String task, LocalDate clearingDate, String market, String taskStatus, Long securityId) { - this.taskTime = taskTime; + public String getTask() { + return task; + } + + public void setTask(String task) { this.task = task; - this.clearingDate = clearingDate; - this.market = market; - this.taskStatus = taskStatus; - this.securityId = securityId; } public LocalTime getTaskTime() { @@ -49,14 +51,6 @@ public class SchedulerNewRequest { this.taskTime = taskTime; } - public String getTask() { - return task; - } - - public void setTask(String task) { - this.task = task; - } - public LocalDate getClearingDate() { return clearingDate; } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java index abdf25d5a..6c46bee95 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/SchedulerUpdateRequest.java @@ -12,30 +12,33 @@ import java.time.LocalDate; import java.time.LocalTime; public class SchedulerUpdateRequest { - @JsonProperty - private String task; + @JsonProperty + public Long id; + @JsonProperty + public String task; @JsonSerialize(using = LocalTimeSerializer.class) @JsonDeserialize(using = LocalTimeDeserializer.class) @JsonProperty - private LocalTime taskTime; - + public LocalTime taskTime; @JsonSerialize(using = LocalDateSerializer.class) @JsonDeserialize(using = LocalDateDeserializer.class) @JsonProperty - private LocalDate clearingDate; - + public LocalDate clearingDate; @JsonProperty - private String market; - + public String market; @JsonProperty - private String taskStatus; - + public String taskStatus; @JsonProperty - private Long securityId; + public Long securityId; - @JsonProperty - private Long id; + public Long getId() { + return id; + } + + public void setId(Long id) { + this.id = id; + } public String getTask() { return task; @@ -84,12 +87,4 @@ public class SchedulerUpdateRequest { public void setSecurityId(Long securityId) { this.securityId = securityId; } - - public Long getId() { - return id; - } - - public void setId(Long id) { - this.id = id; - } } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TimetableNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TimetableNewRequest.java new file mode 100644 index 000000000..ec178490d --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TimetableNewRequest.java @@ -0,0 +1,45 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.schedule; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer; + +import java.time.LocalTime; + +public class TimetableNewRequest { + @JsonProperty + public String task; + + @JsonSerialize(using = LocalTimeSerializer.class) + @JsonDeserialize(using = LocalTimeDeserializer.class) + @JsonProperty + public LocalTime taskTime; + @JsonProperty + public String taskStatus; + + public String getTask() { + return task; + } + + public void setTask(String task) { + this.task = task; + } + + public LocalTime getTaskTime() { + return taskTime; + } + + public void setTaskTime(LocalTime taskTime) { + this.taskTime = taskTime; + } + + public String getTaskStatus() { + return taskStatus; + } + + public void setTaskStatus(String taskStatus) { + this.taskStatus = taskStatus; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TimetableUpdateRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TimetableUpdateRequest.java new file mode 100644 index 000000000..210869118 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TimetableUpdateRequest.java @@ -0,0 +1,56 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.schedule; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalTimeDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSerializer; + +import java.time.LocalTime; + +public class TimetableUpdateRequest { + @JsonProperty + public Long id; + + @JsonProperty + public String task; + + @JsonSerialize(using = LocalTimeSerializer.class) + @JsonDeserialize(using = LocalTimeDeserializer.class) + @JsonProperty + public LocalTime taskTime; + @JsonProperty + public String taskStatus; + + public String getTask() { + return task; + } + + public void setTask(String task) { + this.task = task; + } + + public LocalTime getTaskTime() { + return taskTime; + } + + public void setTaskTime(LocalTime taskTime) { + this.taskTime = taskTime; + } + + public String getTaskStatus() { + return taskStatus; + } + + public void setTaskStatus(String taskStatus) { + this.taskStatus = taskStatus; + } + + public Long getId() { + return id; + } + + public void setId(Long id) { + this.id = id; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TradingCalendarNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TradingCalendarNewRequest.java new file mode 100644 index 000000000..2472babb8 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TradingCalendarNewRequest.java @@ -0,0 +1,44 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.schedule; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; + +import java.time.LocalDate; + +public class TradingCalendarNewRequest { + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty + public LocalDate clearingDate; + @JsonProperty + public Long companyId; + @JsonProperty + public String tradingStatus; + + public LocalDate getClearingDate() { + return clearingDate; + } + + public void setClearingDate(LocalDate clearingDate) { + this.clearingDate = clearingDate; + } + + public Long getCompanyId() { + return companyId; + } + + public void setCompanyId(Long companyId) { + this.companyId = companyId; + } + + public String getTradingStatus() { + return tradingStatus; + } + + public void setTradingStatus(String tradingStatus) { + this.tradingStatus = tradingStatus; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TradingCalendarUpdateRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TradingCalendarUpdateRequest.java new file mode 100644 index 000000000..7fb496a9a --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/TradingCalendarUpdateRequest.java @@ -0,0 +1,55 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.schedule; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; + +import java.time.LocalDate; + +public class TradingCalendarUpdateRequest { + @JsonProperty + public Long id; + + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty + public LocalDate clearingDate; + @JsonProperty + public Long companyId; + @JsonProperty + public String tradingStatus; + + public Long getId() { + return id; + } + + public void setId(Long id) { + this.id = id; + } + + public LocalDate getClearingDate() { + return clearingDate; + } + + public void setClearingDate(LocalDate clearingDate) { + this.clearingDate = clearingDate; + } + + public Long getCompanyId() { + return companyId; + } + + public void setCompanyId(Long companyId) { + this.companyId = companyId; + } + + public String getTradingStatus() { + return tradingStatus; + } + + public void setTradingStatus(String tradingStatus) { + this.tradingStatus = tradingStatus; + } +}