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..c618bad8f --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentInstructionSorter.java @@ -0,0 +1,120 @@ +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"); + } + } +} 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/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..baf2b5080 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 @@ -35,6 +35,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