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())); + } } }