From 9a97006dcbcaeb92346c06f964dd0c5513613db6 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Mon, 18 Sep 2023 17:42:26 +0300 Subject: [PATCH] =?UTF-8?q?account-service=20utility-service=20http://jira?= =?UTF-8?q?.mfd.msk:8088/browse/CLS-527=20DF52=20=D0=BE=D0=B6=D0=B8=D0=B4?= =?UTF-8?q?=D0=B0=D0=BD=D0=B8=D0=B5=20=D0=B0=D0=BA=D1=86=D0=B5=D0=BF=D1=82?= =?UTF-8?q?=D0=B0=20=D0=BF=D0=BE=D0=BB=D1=8C=D0=B7=D0=BE=D0=B2=D0=B0=D1=82?= =?UTF-8?q?=D0=B5=D0=BB=D1=8F=20=D1=87=D0=B5=D1=80=D0=B5=D0=B7=20notificat?= =?UTF-8?q?ion's=D1=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/ClearingAccountService.java | 149 +++++++++++++++++- .../account/service/SDFProcessService.java | 90 +++++++++++ .../utility/service/NotificationService.java | 4 +- .../platform/enumeration/ObjectType.java | 3 +- .../platform/messaging/domain/Consts.java | 1 + 5 files changed, 240 insertions(+), 7 deletions(-) diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java index d5d337440..9a2aff484 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java @@ -15,8 +15,8 @@ import ru.clearing.classes.statics.data.account.ClearingAccount; import ru.clearing.classes.statics.data.company.ClearingMemberCategory; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.relation.Relation; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.clearing.classes.statics.data.sdf.SDf52; -import ru.clearing.platform.dictionary.AbstractDictionary; import ru.clearing.platform.dictionary.ClearingCategoryDictionary; import ru.spcex.clearing.account.errors.AccountError; import ru.spcex.clearing.imdg.IMDGDistributedNames; @@ -28,6 +28,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf0 import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountSdfToStatementRequestPart; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationFeedbackRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; @@ -42,7 +45,6 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; import ru.spcex.platform.utils.collection.Pair; import ru.spcex.platform.utils.enumeration.EnumMessage; -import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IErrorEnumId; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.validation.IValidator; @@ -50,7 +52,6 @@ import ru.spcex.platform.utils.validation.IValidator; import java.time.Instant; import java.util.*; import java.util.function.Function; -import java.util.stream.Collectors; @Service public class ClearingAccountService extends QueueConsumer implements InitializingBean { @@ -70,6 +71,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin private final Imdg relationImdg; private final Imdg clearingMemberCategoryImdg; private final Imdg clearingCategoryImdg; + private final Imdg tradingClearingRegistryImdg; @Autowired public ClearingAccountService(Consumer kafkaQueue, @@ -101,6 +103,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin this.relationImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class); this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCategoryDictionary, ClearingCategoryDictionary.class); + this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); } @Override @@ -114,9 +117,13 @@ public class ClearingAccountService extends QueueConsumer implements Initializin callback(AccountSdf01Request.class) .setFunction(this::accountNewSdf01) .forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put); + callback(StatementRequest.class) .setFunction(this::accountUpdateSdf52) .forDestination(Consts.ACCOUNT_PROCESS_SDF52, callbacks::put); + callback(NotificationFeedbackRequest.class) + .setFunction(this::accountUpdateSdf52_part2Notification) + .forDestination(Consts.ACCOUNT_NOTIFICATION_FEEDBACK, callbacks::put); init(); } @@ -359,7 +366,79 @@ public class ClearingAccountService extends QueueConsumer implements Initializin } } log.debug("Selected to update {} account's", toUpdate.size()); + + // 4. to notification + List notificationIds = new ArrayList<>(); + boolean needWait = false; + { + for (Pair item : toUpdate) { + SDf52 sdf = item.getFirst(); + Account account = item.getSecond(); + AccountStatus newStatus = sdfProcessService.parseSdf52Status(sdf.getStatus()); + if (newStatus == null) { + throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus()); + } + if (newStatus.equalsByKey(account.getStatus())) { + // одинаковых обычно не бывает. + continue; + } + if (AccountStatus.BLOCKED == newStatus || AccountStatus.CLOSE == newStatus) { // статус 0/2 + notificationIds.add(sendNotificationRequest(ObjectType.account_block, account, newStatus)); + needWait = true; + } + if (AccountStatus.ACTIVE == newStatus) { // статус 1/3 + notificationIds.add(sendNotificationRequest(ObjectType.account_active, account, newStatus)); + needWait = true; + } + } + } + if (needWait) { + SDFProcessService.SDF52WaitingData data = new SDFProcessService.SDF52WaitingData(systemRequest, toProcessSDF53, toUpdate, notificationIds); + sdfProcessService.putNotificationWaiting(data); + log.info("Account's need wait accept over notification. Waiting user. GroupId={}", groupId); + } else { + log.info("No one account for wait notification reply. Do immediatly stage 2. groupId={}", groupId); + accountUpdateSdf52_part2(systemRequest, toProcessSDF53, toUpdate, null); + } + + log.debug("successfully processed, grouping id={} with {} accounts.", + groupId, toUpdate.size()); + return null; + } + + protected RequestInfoUpdate accountUpdateSdf52_part2Notification(BaseRequest secondSystemRequest2) { + log.debug("accountUpdateSdf52 NotificationFeedbackRequest received, id={}", secondSystemRequest2.getId()); + SDFProcessService.SDF52WaitingData trigger = sdfProcessService.onNotificationResponse(secondSystemRequest2.getRequestPayload()); + if (trigger == null) { + log.trace("Not triggered, continue waiting"); + return null; + } else { + log.debug("Triggered, groupId={}, notification's status: {}, last notification user id={}", + trigger.getGroupId(), trigger.globalStatus, secondSystemRequest2.getUserId()); + if (NotificationStatus.ACPT.equalsByKey(trigger.globalStatus)) { + return accountUpdateSdf52_part2(trigger.systemRequest, trigger.toProcessSDF53, trigger.toUpdate, secondSystemRequest2); + } else if (NotificationStatus.CNCL.equalsByKey(trigger.globalStatus)) { + log.debug("groupId={} rejected by user {}.", trigger.getGroupId(), secondSystemRequest2.getUserId()); + return null; + } else { + log.warn("Unknown waiting group status \"{}\"", trigger.globalStatus); + return null; + } + } + } + + protected RequestInfoUpdate accountUpdateSdf52_part2(BaseRequest systemRequest1, + List> toProcessSDF53, + List> toUpdate, + BaseRequest secondSystemRequest2) { + // 5. from notification: + StatementRequest req = systemRequest1.getRequestPayload(); + Long groupId = req.getGroupId(); + log.info("Continue Sdf52 groupId={}, first request id={}, second request id={}", + groupId, systemRequest1 == null ? null : systemRequest1.getId(), secondSystemRequest2 == null ? null : secondSystemRequest2.getId()); + int countOfUpdated = 0; + List toTCRRequests = new ArrayList<>(); ImdgTransaction imdgTransaction = imdgProvider.newTransaction(); boolean txOk = false; imdgTransaction.beginTransaction(); @@ -374,14 +453,28 @@ public class ClearingAccountService extends QueueConsumer implements Initializin if (newStatus == null) { throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus()); } + // За долгое время ожидания пользователя счёт мог быть обновлён, прочитать его ещё раз + account = accountImdg.getSingleObjectByID(account.getId()); + item.setSecond(account); // В toProcessSDF53 есть тоже account, но он на следующих этапах не используется toProcessSDF53.add(new MutableTriple<>(sdf, account, SDFProcessService.SDF_STATUS_OK)); if (!newStatus.equalsByKey(account.getStatus())) { // Обновление счёта account.setStatus(newStatus.getKey()); account.setUpdated(now); accountImdg.update(account); - log.trace("S_DF52[{}] do update status to {} for account[{}]", sdf.getId(), newStatus.getKey(), account.getId()); + log.trace("transaction {}, S_DF52[{}] do update status to {} for account[{}]", + imdgTransaction, sdf.getId(), newStatus.getKey(), account.getId()); + if (AccountStatus.ACTIVE == newStatus) { + TradingClearingRegistryUpdateRequest tcrReq = new TradingClearingRegistryUpdateRequest(); + tcrReq.setId(account.getId()); + tcrReq.setCompanyId(account.getCompanyId()); + // tcrReq.setId(); // select TCR; - выполнить поиск вне этой транзакции, чтобы + tcrReq.setStatus(account.getStatus()); + toTCRRequests.add(tcrReq); + } countOfUpdated++; + } else { + log.info("Account[{}] \"{}\" do not updated - same status \"{}\"", account.getId(), account.getAccount(), newStatus); } } txOk = true; @@ -389,12 +482,28 @@ public class ClearingAccountService extends QueueConsumer implements Initializin if (txOk) { imdgTransaction.commitTransaction(); } else { - log.debug("failed update, clearing accounts, rollback transaction. Request id={}; sdf52 groupId={}", systemRequest.getId(), groupId); + log.debug("failed update, clearing accounts, rollback transaction {}. Request id={}; sdf52 groupId={}", + imdgTransaction, systemRequest1 == null ? null : systemRequest1.getId(), groupId); imdgTransaction.rollbackTransaction(); } } if (txOk) { + for (TradingClearingRegistryUpdateRequest tcrReq : toTCRRequests) { + TradingClearingRegistry tcr = tradingClearingRegistryImdg.getFirstObjectByFieldValues(Map.of( + "moneyAccountId", tcrReq.getMoneyAccountId(), + "companyId", tcrReq.getCompanyId())); + if (tcr == null) { + log.warn("TradingClearingRegistry by moneyAccountId={} and companyId={} not found, do not send update TCR.", + tcrReq.getMoneyAccountId(), tcrReq.getCompanyId()); + } else { + tcrReq.setId(tcr.getId()); + String destination = Consts.DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE; + log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(tcrReq)); + kafkaSender.sendRequestToQueue(destination, tcrReq); + } + } + sdfProcessService.process(req, toProcessSDF53); } @@ -422,4 +531,34 @@ public class ClearingAccountService extends QueueConsumer implements Initializin log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request)); kafkaSender.sendRequestToQueue(destination, request); } + + public void sendTCRUpdateRequest(Long groupingSdf01Id, Long groupSdf02Id, List results) { + final String destination = Consts.STATEMENT_PROCESS; + StatementRequest request = new StatementRequest(); + request.setGroupId(groupingSdf01Id); + request.setChildGenerationId(groupSdf02Id); + request.setAccountCreationResults(results); + request.setContinueSdf(true); + request.setTable(SdfTable.SDF_01); // по нему запрос получили + log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request)); + kafkaSender.sendRequestToQueue(destination, request); + } + + public Long sendNotificationRequest(ObjectType objectType, Account account, AccountStatus requestToStatus) { + String statusText = AccountStatus.BLOCKED == requestToStatus ? "Заблокирован" + : AccountStatus.CLOSE == requestToStatus ? "Закрыт" + : AccountStatus.ACTIVE == requestToStatus ? "Разблокирован" : + requestToStatus.getKey(); + String message = String.format("Для счета «%s» будет изменен статус на «%s»", account.getAccount(), statusText); + final String destination = Consts.NOTIFICATION_NEW; + NotificationNewRequest request = new NotificationNewRequest(); + request.setObjectType(objectType.getKey()); + request.setObjectId(account.getId()); + request.setPriority(Priority.HIGH.getKey()); + request.setComment(message); + log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request)); + Long rKey = kafkaSender.sendRequestToQueue(destination, request); + log.trace("About account.id={} send notification request id={}", account.getId(), rKey); + return rKey; + } } diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/SDFProcessService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/SDFProcessService.java index 6f5abc994..3321cd9ac 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/SDFProcessService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/SDFProcessService.java @@ -9,18 +9,23 @@ import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.sdf.SDf52; import ru.clearing.classes.statics.data.sdf.SDf53; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationFeedbackRequest; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.enumeration.AccountStatus; +import ru.spcex.platform.enumeration.NotificationStatus; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgId; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.ImdgTransaction; +import ru.spcex.platform.utils.collection.Pair; import java.time.Instant; import java.util.*; +import java.util.concurrent.CopyOnWriteArrayList; @Service public class SDFProcessService { @@ -144,4 +149,89 @@ public class SDFProcessService { log.debug("Send ExportToFileRequest({}, {}) message id={} to kafka \"{}\"", groupId, request.getNameOfTable(), msgId, destination); } + + // --------- notification apply system ----------- + public static class SDF52WaitingData { + public BaseRequest systemRequest; + public List> toProcessSDF53; + public List> toUpdate; + public HashSet notAnsweredNotifications = new HashSet<>(); + public String globalStatus; + + public SDF52WaitingData(BaseRequest systemRequest, List> toProcessSDF53, List> toUpdate, + Collection notAnsweredNotifications) { + this.systemRequest = systemRequest; + this.toProcessSDF53 = toProcessSDF53; + this.toUpdate = toUpdate; + this.notAnsweredNotifications = new HashSet<>(notAnsweredNotifications); + } + + public Long getGroupId() { + if (systemRequest != null && systemRequest.getRequestPayload() != null) + return systemRequest.getRequestPayload().getGroupId(); + return null; + } + } + + protected CopyOnWriteArrayList waitingList = new CopyOnWriteArrayList<>(); + + public boolean putNotificationWaiting(SDF52WaitingData data) { + for (SDF52WaitingData i : waitingList) { + if (Objects.equals(i.getGroupId(), data.getGroupId())) { + log.warn("groupId={} already waiting", data.getGroupId()); + } + } + data.globalStatus = null; + waitingList.add(data); + return true; + } + + /** + * @return triggered SDF52WaitingData or null + */ + public SDF52WaitingData onNotificationResponse(NotificationFeedbackRequest onNotification) { + final Long nId = onNotification.getNotificationId(); + if (nId == null) + return null; + SDF52WaitingData inWaitingLst = null; + for (SDF52WaitingData i : waitingList) + synchronized (i) { + if (i.notAnsweredNotifications.contains(nId)) { + inWaitingLst = i; + break; + } + } + if (inWaitingLst == null) { + log.debug("SDF52 waiting list not found for notificationId={}", nId); + return null; + } + boolean returnTrigger = false; + synchronized (inWaitingLst) { + boolean removedId = inWaitingLst.notAnsweredNotifications.remove(nId); + if (!removedId) { // never + log.warn("Can not remove notificationId={} from waiting list", nId); + } + if (NotificationStatus.ACPT.equalsByKey(onNotification.getNotificationStatus())) { + if (inWaitingLst.notAnsweredNotifications.isEmpty()) { + log.trace("Waiting list is empty, by notificationId={}", nId); + inWaitingLst.globalStatus = onNotification.getNotificationStatus(); + returnTrigger = true; + } + } else if (NotificationStatus.CNCL.equalsByKey(onNotification.getNotificationStatus())) { + log.trace("Waiting list cancel, by notificationId={}", nId); + inWaitingLst.globalStatus = onNotification.getNotificationStatus(); + returnTrigger = true; + } else { + log.warn("Unknown notification[{}] status={}", onNotification.getNotificationId(), onNotification.getNotificationStatus()); + } + if (returnTrigger || inWaitingLst.notAnsweredNotifications.isEmpty()) { + log.debug("Remove waiting list by notificationId={} (list groupId={})", nId, inWaitingLst.getGroupId()); + waitingList.remove(inWaitingLst); + } + } + if (returnTrigger) + return inWaitingLst; + else + return null; + } } 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 index d47b0b6d6..9585acf82 100644 --- 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 @@ -23,7 +23,7 @@ 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; +import static ru.spcex.platform.enumeration.ObjectType.*; @Service public class NotificationService extends QueueConsumer implements InitializingBean { @@ -94,6 +94,8 @@ public class NotificationService extends QueueConsumer implements InitializingBe String objectType = notification.getObjectType(); if (objectType.equalsIgnoreCase(statement.getKey())) { kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(notification.getId(), notification.getNotificationStatus())); + } else if (objectType.equalsIgnoreCase(account_active.getKey()) || objectType.equalsIgnoreCase(account_block.getKey())) { + kafkaSender.sendRequestToQueue(ACCOUNT_NOTIFICATION_FEEDBACK, buildFeedbackRequest(notification.getId(), notification.getNotificationStatus())); } else { log.warn("Unsupported notification type {}", objectType); } diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ObjectType.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ObjectType.java index 56a090cc6..d3f1475c0 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ObjectType.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/ObjectType.java @@ -3,7 +3,8 @@ package ru.spcex.platform.enumeration; import ru.spcex.platform.utils.enumeration.IEnumKey; public enum ObjectType implements IEnumKey { - statement("STMT"), vfrs("VFRS"), rgst("RGST"), gateway("GTWY"), session("SESN"); + statement("STMT"), vfrs("VFRS"), rgst("RGST"), gateway("GTWY"), session("SESN"), + account_block("ACCB"), account_active("ACCA"); private final String 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 379136ac0..40d1fb693 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 @@ -111,6 +111,7 @@ public interface Consts { String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block"; String CLEARING_NOTIFICATION_FEEDBACK = "clearing-notification-feedback"; + String ACCOUNT_NOTIFICATION_FEEDBACK = "account-notification-feedback"; String NOTIFICATION_NEW = "notification-new"; String NOTIFICATION_UPDATE = "notification-update";