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 index 1103e95a6..887fda87d 100644 --- 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 @@ -6,118 +6,33 @@ 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; + private final SdfCreatorBySTLDPayment sdfCreator; + private final PaymentUpdateBySdf04 paymentUpdater; @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; + public ClearingService(SdfCreatorBySTLDPayment sdfCreator, PaymentUpdateBySdf04 paymentUpdater) { + this.sdfCreator = sdfCreator; + this.paymentUpdater = paymentUpdater; this.executor = Executors.newSingleThreadExecutor(); } @Scheduled(cron = "${clearing-service.scheduler.check-payment-instruction}") - public void run() { - executor.execute(this::createSdfFromPaymentInstructionSTLD); + public void sdfCreate() { + log.info("creating sdf03/11 from STLD payments task added to queue"); + executor.execute(sdfCreator::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); - } - } + public void paymentUpdateBySdf04(Long sdf04GroupId) { + log.info("adding sdf03/11 from STLD payments task to queue"); + executor.execute(() -> paymentUpdater.updatePayments(sdf04GroupId)); } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentUpdateBySdf04.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentUpdateBySdf04.java new file mode 100644 index 000000000..458528d20 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentUpdateBySdf04.java @@ -0,0 +1,47 @@ +package ru.spcex.clearing.service; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.payment.PaymentInstruction; +import ru.clearing.classes.statics.data.sdf.SDf04; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.enumeration.TransactionStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.Collection; +import java.util.Map; + +@Component +public class PaymentUpdateBySdf04 { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg paymentImdgs; + private final Imdg sdf04Imdg; + + public PaymentUpdateBySdf04(ImdgProvider imdgProvider) { + this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class); + this.sdf04Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class); + } + + public void updatePayments(Long sdf04GroupId) { + Collection sDf04s = sdf04Imdg.getCollectionObjectsByFieldValues(Map.of("generationId", sdf04GroupId)); + for (SDf04 sDf04 : sDf04s) { + String docnm_ref = sDf04.getDocnm_ref(); + long paymentId; + try { + paymentId = Long.parseLong(docnm_ref); + } catch (NumberFormatException e) { + log.warn("sdf04.id={} cannot parse docnm_ref={} as payment id", sDf04.getId(), sDf04.getDocnm_ref()); + continue; + } + PaymentInstruction payment = paymentImdgs.getSingleObjectByFieldValues(Map.of("id", paymentId)); + if ("OK!".equals(sDf04.getImp_result())) { + payment.setTransactionStatus(TransactionStatus.ok.getKey()); + } else { + payment.setTransactionStatus(TransactionStatus.fail.getKey()); + } + paymentImdgs.update(payment); + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf04Receiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf04Receiver.java new file mode 100644 index 000000000..08db5dde1 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf04Receiver.java @@ -0,0 +1,28 @@ +package ru.spcex.clearing.service; + +import org.apache.kafka.clients.consumer.Consumer; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; + +@Service +public class Sdf04Receiver extends QueueConsumer implements InitializingBean { + private final ClearingService clearingService; + public Sdf04Receiver(Consumer kafkaQueue, ClearingService clearingService) { + super(kafkaQueue); + this.clearingService = clearingService; + } + + @Override + public void afterPropertiesSet() { + callback(Sdf04Request.class) + .setConsumer(event -> { + Sdf04Request requestPayload = event.getRequestPayload(); + clearingService.paymentUpdateBySdf04(requestPayload.getGroupId()); + }) + .forDestination(Consts.SDF04_PROCESS, callbacks::put); + init(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java new file mode 100644 index 000000000..5ff6e7f24 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/SdfCreatorBySTLDPayment.java @@ -0,0 +1,112 @@ +package ru.spcex.clearing.service; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; +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.clearing.service.order.PaymentBatchInfo; +import ru.spcex.clearing.service.order.PaymentInstructionSorter; +import ru.spcex.clearing.service.order.Sdf03Builder; +import ru.spcex.clearing.service.order.Sdf11Builder; +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.stream.Collectors; + +@Component +public class SdfCreatorBySTLDPayment { + private final Imdg paymentImdgs; + private final ImdgId idGenerator; + private final PaymentInstructionSorter senderGroupSorter; + private final Imdg sdf03Imdg; + private final Imdg sdf11Imdg; + private final KafkaSender kafkaSender; + + @Autowired + public SdfCreatorBySTLDPayment(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; + } + + public 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/order/PaymentBatchInfo.java similarity index 98% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentBatchInfo.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/order/PaymentBatchInfo.java index fff82f55c..4f96a3856 100644 --- 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/order/PaymentBatchInfo.java @@ -1,4 +1,4 @@ -package ru.spcex.clearing.service; +package ru.spcex.clearing.service.order; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.spcex.platform.enumeration.ClearingMemberCategoryD; 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/order/PaymentInstructionSorter.java similarity index 97% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/PaymentInstructionSorter.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/order/PaymentInstructionSorter.java index f5f59d15e..1d8c701ab 100644 --- 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/order/PaymentInstructionSorter.java @@ -1,4 +1,4 @@ -package ru.spcex.clearing.service; +package ru.spcex.clearing.service.order; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,7 +40,7 @@ public class PaymentInstructionSorter { this.messageResolver = messageResolver; } - PaymentBatchInfo sortCompanyPayments(Long generationId, Long senderId, List payments) { + public 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)); 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/order/Sdf03Builder.java similarity index 96% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf03Builder.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/order/Sdf03Builder.java index fc32df0e2..a16f02452 100644 --- 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/order/Sdf03Builder.java @@ -1,7 +1,8 @@ -package ru.spcex.clearing.service; +package ru.spcex.clearing.service.order; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.clearing.classes.statics.data.sdf.SDf03; +import ru.spcex.clearing.service.SpecifUtil; import ru.spcex.platform.enumeration.Sender; import ru.spcex.platform.utils.time.TimeUtil; 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/order/Sdf11Builder.java similarity index 96% rename from clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/Sdf11Builder.java rename to clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/order/Sdf11Builder.java index 2d37e878b..310a77b1d 100644 --- 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/order/Sdf11Builder.java @@ -1,7 +1,8 @@ -package ru.spcex.clearing.service; +package ru.spcex.clearing.service.order; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.clearing.classes.statics.data.sdf.SDf11; +import ru.spcex.clearing.service.SpecifUtil; import ru.spcex.platform.enumeration.Sender; import ru.spcex.platform.utils.time.TimeUtil; 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 index 2d97429b9..554033e55 100644 --- 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 @@ -3,7 +3,7 @@ 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"); + stld("STLD"), notSent("NSNT"), cher("CHER"), sent("SENT"), ok("OK"), fail("FAIL"); private final String key; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index 0e11d3bb8..1ecdd1625 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -50,6 +50,12 @@ public class QueueConsumer implements AutoCloseable { this.supportStartOffsetTimeWindow = false; } + /** + * используем этот конструктор, если хотим класть + * в кафку "ответ" - информацию о статусе обработки команд + * @param kafkaQueue + * @param kafkaResponseQueue + */ public QueueConsumer(Consumer kafkaQueue, Producer kafkaResponseQueue) { this(kafkaQueue); this.producer = kafkaResponseQueue;