From 9c3e8737a079a0ccf8b97c3d113d6eeb0d612ca9 Mon Sep 17 00:00:00 2001 From: aalehin Date: Mon, 5 Jun 2023 17:42:02 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-339 --- NotificationService set up --- .../balance/service/Sdf16Executor.java | 3 +- .../clearing/utility/config/KafkaConfig.java | 30 ++++++ .../utility/service/NotificationService.java | 94 +++++++++++++++++++ .../enumeration/NotificationStatus.java | 21 +++++ .../platform/messaging/domain/Consts.java | 7 +- .../NotificationFeedbackRequest.java | 17 ++++ .../utilities/NotificationUpdateRequest.java | 27 ++++++ 7 files changed, 197 insertions(+), 2 deletions(-) create mode 100644 clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java create mode 100644 platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/NotificationStatus.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationUpdateRequest.java diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf16Executor.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf16Executor.java index eb9080003..00df68179 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf16Executor.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf16Executor.java @@ -51,7 +51,8 @@ public class Sdf16Executor extends AbstractExecutor { LoggingService errorLogger, ImdgProvider imdgProvider, AccountBalanceService accountBalanceService, - IMessageResolver errorResolver, KafkaSender kafaSender) { + IMessageResolver errorResolver, + KafkaSender kafaSender) { this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); this.sdf17Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf17, SDf17.class); this.accountBalanceImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class); diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java index 2dde7ea72..bb9abc2af 100644 --- a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/config/KafkaConfig.java @@ -3,13 +3,22 @@ package ru.spcex.clearing.utility.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.annotation.Qualifier; 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 org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +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.clearing.utility.config.settings.UtilityServiceSettings; +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 { @@ -26,4 +35,25 @@ public class KafkaConfig { return KafkaProducerFactory.producer(settings.getKafkaProducer()); } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + + @Autowired + @Bean + public KafkaSender kafkaSender(@Qualifier("kafkaTemplate") KafkaTemplate kafkaTemplate, + ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); + } } diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java new file mode 100644 index 000000000..cb5b13b60 --- /dev/null +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/NotificationService.java @@ -0,0 +1,94 @@ +package ru.spcex.clearing.utility.service; + +import org.apache.kafka.clients.consumer.Consumer; +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.misc.Notification; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.*; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.time.Instant; +import java.time.LocalDate; + +import static ru.spcex.clearing.platform.messaging.domain.Consts.*; +import static ru.spcex.platform.enumeration.NotificationStatus.PEND; +import static ru.spcex.platform.enumeration.ObjectType.statement; + +@Service +public class NotificationService extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg notificationMap; + private final KafkaSender kafkaSender; + + @Autowired + public NotificationService(Consumer kafkaQueue, + ImdgProvider imdgProvider, + KafkaSender kafkaSender) { + super(kafkaQueue); + this.notificationMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Notification, Notification.class); + this.kafkaSender = kafkaSender; + } + + @Override + public void afterPropertiesSet() { + callback(NotificationNewRequest.class) + .setConsumer(this::notificationNewRequest) + .forDestination(NOTIFICATION_NEW, callbacks::put); + callback(NotificationUpdateRequest.class) + .setConsumer(this::notificationUpdateRequest) + .forDestination(NOTIFICATION_UPDATE, callbacks::put); + init(); + } + + private void notificationNewRequest(BaseRequest notificationNewRequestBaseRequest) { + log.info("Starting notificationNewRequest processing..."); + NotificationNewRequest request = notificationNewRequestBaseRequest.getRequestPayload(); + Notification notification = notificationBuilderFromNotificationNewRequest(request); + notificationMap.insert(notification); + log.info("NotificationNewRequest successfully processed! Generated id:\t{}", notification.getId()); + } + + private void notificationUpdateRequest(BaseRequest notificationUpdateRequestBaseRequest) { + log.info("Starting notificationUpdateRequest processing..."); + NotificationUpdateRequest request = notificationUpdateRequestBaseRequest.getRequestPayload(); + Notification notification = notificationMap.getSingleObjectByID(request.getId()); + notification.setNotificationStatus(request.getNotificationStatus()); + notification.setUpdated(Instant.now()); + notificationMap.update(notification); + sendFeedbackString(notification.getObjectType(), notification.getNotificationStatus()); + log.info("NotificationUpdateRequest successfully processed!"); + } + + private Notification notificationBuilderFromNotificationNewRequest(NotificationNewRequest request) { + Notification notification = new Notification(); + notification.setClearingDate(LocalDate.now()); + notification.setSenderId(request.getSenderId()); + notification.setAddresseeId(request.getAddresseeId()); + notification.setObjectType(request.getObjectType()); + notification.setObjectId(request.getObjectId()); + notification.setNotificationStatus(PEND.getKey()); + notification.setCreated(Instant.now()); + notification.setUpdated(Instant.now()); + return notification; + } + + private void sendFeedbackString(String objectType, String status){ + if (objectType.equalsIgnoreCase(statement.getKey())){ + kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(status)); + } + } + + private NotificationFeedbackRequest buildFeedbackRequest(String status){ + NotificationFeedbackRequest request = new NotificationFeedbackRequest(); + request.setNotificationStatus(status); + return request; + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/NotificationStatus.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/NotificationStatus.java new file mode 100644 index 000000000..553981867 --- /dev/null +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/NotificationStatus.java @@ -0,0 +1,21 @@ +package ru.spcex.platform.enumeration; + +import ru.spcex.platform.utils.enumeration.IEnumKey; + +public enum NotificationStatus implements IEnumKey { + PEND("PEND"), + CNCL("CNCL"), + ACPT("ACPT"), + READ("READ"); + + private final String key; + + NotificationStatus(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 7101cbf0f..717696611 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 @@ -102,6 +102,11 @@ public interface Consts { String DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE = "trading-clearing-registry-update"; String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block"; + String CLEARING_NOTIFICATION_FEEDBACK = "clearing-notification-feedback"; + String NOTIFICATION_NEW = "notification-new"; + String NOTIFICATION_UPDATE = "notification-update"; + + String DESTINATION_SDF08_NEW = "s-df-08-new"; String DESTINATION_SDF02_NEW = "s-df-02-new"; @@ -137,7 +142,7 @@ public interface Consts { String BALANCE_ACCOUNT_UPDATE = "balance-account-update"; String CONTINUE_CLEARING = "continue-clearing"; String LAUNCHER_NEW = "launcher-new"; - String NOTIFICATION_NEW = "notification-new"; + String REQUEST_INFO_UPDATE = "request-info-update"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java new file mode 100644 index 000000000..f58461bfc --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationFeedbackRequest.java @@ -0,0 +1,17 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.utilities; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; + +public class NotificationFeedbackRequest { + @JsonProperty + public String notificationStatus; + + public String getNotificationStatus() { + return notificationStatus; + } + + public void setNotificationStatus(String notificationStatus) { + this.notificationStatus = notificationStatus; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationUpdateRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationUpdateRequest.java new file mode 100644 index 000000000..3a14578d3 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/utilities/NotificationUpdateRequest.java @@ -0,0 +1,27 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.utilities; + +import com.fasterxml.jackson.annotation.JsonProperty; + +public class NotificationUpdateRequest { + @JsonProperty + public Long id; + + @JsonProperty + public String notificationStatus; + + public Long getId() { + return id; + } + + public void setId(Long id) { + this.id = id; + } + + public String getNotificationStatus() { + return notificationStatus; + } + + public void setNotificationStatus(String notificationStatus) { + this.notificationStatus = notificationStatus; + } +}