---
NotificationService set up
This commit is contained in:
aalehin 2023-06-05 17:42:02 +03:00
parent 415d5697d8
commit 9c3e8737a0
7 changed files with 197 additions and 2 deletions

View file

@ -51,7 +51,8 @@ public class Sdf16Executor extends AbstractExecutor<SDf16> {
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);

View file

@ -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<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Autowired
@Bean
public KafkaSender kafkaSender(@Qualifier("kafkaTemplate") KafkaTemplate<String, Object> kafkaTemplate,
ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
return KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
.imdgProvider(s -> {
Imdg<RequestInfo> imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class);
return imdg::insert;
})
.build();
}
}

View file

@ -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<Notification> notificationMap;
private final KafkaSender kafkaSender;
@Autowired
public NotificationService(Consumer<String, Object> 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<NotificationNewRequest> 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<NotificationUpdateRequest> 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;
}
}

View file

@ -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;
}
}

View file

@ -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";

View file

@ -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;
}
}

View file

@ -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;
}
}