From 848c26d69aa1be6d79c0755d1fc81bfec5d53958 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Thu, 17 Aug 2023 17:49:44 +0300 Subject: [PATCH] =?UTF-8?q?scheduler-service=20OUTV=20=D1=88=D0=B5=D0=B4?= =?UTF-8?q?=D1=83=D0=BB=D0=B5=D1=80=20=D1=80=D0=B5=D0=B4=D0=B8=D1=80=D0=B5?= =?UTF-8?q?=D0=BA=D1=82=D0=B8=D1=82=20=D0=B5=D0=B3=D0=BE=20=D0=B2=20=D0=BE?= =?UTF-8?q?=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D1=8C=20PAYMENT=5FINSTRUCTION=5F?= =?UTF-8?q?CLEARING=5FOUTBOUND=5FACTION?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduler/service/LauncherService.java | 40 ++++++++++++++++++- 1 file changed, 39 insertions(+), 1 deletion(-) diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherService.java index 65d2d190d..8f6aa7b79 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherService.java @@ -3,6 +3,7 @@ package ru.spcex.clearing.scheduler.service; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -11,8 +12,10 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.scheduler.Launcher; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.ActionType; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; @@ -21,8 +24,11 @@ import ru.spcex.clearing.util.security.UserRoleVerification; import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.log.ExceptionUtils; import java.time.Instant; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; import static ru.spcex.platform.utils.enumeration.IEnumKey.getEnumByKey; @@ -74,7 +80,39 @@ public class LauncherService extends QueueConsumer implements InitializingBean { launcher.setUpdated(created); launcherMap.insert(launcher); log.debug("successfully processed, new id {}", launcher.getId()); - kafkaProducer.send(new ProducerRecord<>("launcher-" + launcher.getTask(), userRequest)); + if (Task.dbfExport_OUTV.equalsByKey(req.getTaskName())) { + processRequestOUTV(Task.dbfExport_OUTV, req); + } else { + kafkaProducer.send(new ProducerRecord<>("launcher-" + launcher.getTask(), userRequest)); + } + } + + private void processRequestOUTV(Task dbfExport__OUTV, LauncherCommandRequest launcherNew) { + log.debug("Redirect {} to queue {}", dbfExport__OUTV, Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION); + var req = new PIClearingOutbondActionNewRequest(); + req.setSenderId(launcherNew.getSenderId()); + req.setAddresseeId(launcherNew.getAddresseeId()); + req.setPaymentPurpose(launcherNew.getPaymentPurpose()); + req.setCreditLeg_amount(launcherNew.getCreditLeg_amount()); + req.setCreditLeg_accountId(launcherNew.getCreditLeg_accountId()); + req.setDebitLeg_accountId(launcherNew.getDebitLeg_accountId()); + + final String destination = Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION; + BaseRequest request = new BaseRequest<>(); + request.setId(null); + request.setActionType(ActionType.NEW); + request.setRequestPayload(req); + request.setUserId(launcherNew.getUserId()); + Future send = kafkaProducer.send(new ProducerRecord<>(destination, request)); + try { + send.get(); + log.debug("Request id={} by user {} send to queue \"{}\". Action: {}", request.getId(), request.getUserId(), destination, request.getActionType()); + } catch (InterruptedException e) { + log.error("Request id={} by user {} send to queue \"{}\". Action: {}; error send: thread interrupted", request.getId(), request.getUserId(), destination, request.getActionType()); + Thread.currentThread().interrupt(); + } catch (ExecutionException e) { + log.error("Request id={} by user {} send to queue \"{}\". Action: {}; error send: {}", request.getId(), request.getUserId(), destination, request.getActionType(), ExceptionUtils.getStackTrace(e.getCause())); + } } }