From e47ff426021a4e6c5e593e06378fab469878a35f Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 11 Jul 2025 14:12:28 +0300 Subject: [PATCH] Obligation Admission PAUSE notification --- .../notification/NotificationSender.java | 6 +- .../clearing/session/state/DataEnum.java | 1 + .../action/ObligationAdmissionAction.java | 4 ++ ...eWorkflowStatusWithNotificationAction.java | 59 +++++++++++++++++++ .../action/WorkflowStatusConfiguration.java | 7 ++- .../clearing/session/state/model/RgsErr.java | 29 +++++++++ 6 files changed, 103 insertions(+), 3 deletions(-) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/PauseWorkflowStatusWithNotificationAction.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/model/RgsErr.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/notification/NotificationSender.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/notification/NotificationSender.java index 198b79924..0571fdeed 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/notification/NotificationSender.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/notification/NotificationSender.java @@ -21,12 +21,16 @@ public class NotificationSender { } public void sendNotification(ObjectType objType, String comment, Priority priority) { + sendNotification(objType, comment, priority, null); + } + + public void sendNotification(ObjectType objType, String comment, Priority priority, Long objectId) { NotificationNewRequest reviseNotification = new NotificationNewRequest(); reviseNotification.setObjectType(objType.getKey()); reviseNotification.setComment(comment); reviseNotification.setPriority(priority.getKey()); + reviseNotification.setObjectId(objectId); kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, reviseNotification); log.info("notification {} sent {}", objType, comment); } - } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java index 47a888ac0..e75092e32 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/DataEnum.java @@ -11,6 +11,7 @@ public enum DataEnum { dealsPreparedCurrency, counterPartyId, //Long obligationAdmissionStashedRgs, //Map + erroneousRegistries, //List paymentInstructionReturnMkr, //List paymentInstructionSecurity, //List paymentsWereCreatedDepositReturn, //Collection diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java index 7ee19d66a..c3c585c6c 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/ObligationAdmissionAction.java @@ -28,6 +28,7 @@ import ru.spcex.clearing.service.validation.admission.TcrActiveValidationRule; import ru.spcex.clearing.session.stage.TaskType; import ru.spcex.clearing.session.state.DataEnum; import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.clearing.session.state.model.RgsErr; import ru.spcex.platform.enumeration.RegistryStatus; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @@ -78,6 +79,7 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa .collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket()))); log.info("found {} ({} groups) registries with sessionId {}", registries.size(), byGroups.size(), sessionId); Map rgsToUpdate = new HashMap<>(); + List errs = new ArrayList<>(); for (Map.Entry> grpEntry : byGroups.entrySet()) { EnumMessage groupError = null; List rgsGroup = grpEntry.getValue(); @@ -86,6 +88,7 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa Optional error = validator.tillFirstError(); if (error.isPresent()) { groupError = error.get(); + errs.add(new RgsErr(rgs, messageResolver.resolve(groupError))); break; } } @@ -100,6 +103,7 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa } ctx.getExtendedState().getVariables().put(DataEnum.obligationAdmissionStashedRgs, rgsToUpdate); + ctx.getExtendedState().getVariables().put(DataEnum.erroneousRegistries, errs); // if (wasNack) { // } else { // registryImdg.putAll(rgsToUpdate); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/PauseWorkflowStatusWithNotificationAction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/PauseWorkflowStatusWithNotificationAction.java new file mode 100644 index 000000000..8029606a3 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/PauseWorkflowStatusWithNotificationAction.java @@ -0,0 +1,59 @@ +package ru.spcex.clearing.session.state.action; + +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; +import org.springframework.statemachine.StateContext; +import org.springframework.statemachine.state.State; +import ru.spcex.clearing.notification.NotificationSender; +import ru.spcex.clearing.session.stage.TaskType; +import ru.spcex.clearing.session.state.DataEnum; +import ru.spcex.clearing.session.state.SsnEvent; +import ru.spcex.clearing.session.state.model.RgsErr; +import ru.spcex.platform.enumeration.ObjectType; +import ru.spcex.platform.enumeration.Priority; +import ru.spcex.platform.enumeration.WorkflowStatus; +import ru.spcex.platform.imdg.api.ImdgProvider; + +public class PauseWorkflowStatusWithNotificationAction extends UpdateWorkflowStatusAction { + private final NotificationSender notificationSender; + + public PauseWorkflowStatusWithNotificationAction(ImdgProvider imdgProvider, + NotificationSender notificationSender) { + super(imdgProvider, WorkflowStatus.Pause); + this.notificationSender = notificationSender; + } + + @Override + protected void actualExecute(StateContext ctx) { + super.actualExecute(ctx); + @SuppressWarnings("unchecked") + List errs = ctx.getExtendedState().get(DataEnum.erroneousRegistries, List.class); + errs = errs == null ? Collections.emptyList() : errs; + String comment = errs.stream() + .limit(100) + .map(rgsWr -> "ТКР %s: %s".formatted(rgsWr.getRgs().getTradingClearingRegistry(), rgsWr.getErr())) + .collect(Collectors.joining(";\r\n", notificationPrefix(ctx), "")); + Long ssnId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class); + notificationSender.sendNotification(ObjectType.session, comment, Priority.HIGH, ssnId); + ctx.getExtendedState().getVariables().remove(DataEnum.erroneousRegistries); + } + + private String notificationPrefix(StateContext ctx) { + State target = ctx.getTarget(); + if (target == null) return ""; + TaskType id = target.getId(); + if (id == null) return ""; + switch (id) { + case PSEUDO_waitAfterObligationAdmissionError -> { + return "Ряд обязательств не допущены к клирингу!!!\r\n"; + } +// case PSEUDO_waitAfterInspectionError -> { +// return "Ряд обязательств подлежит исключению из клирингового пула!!!\r\n"; +// } + default -> { + return "Ошибка клиринговой сессии\r\n"; + } + } + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/WorkflowStatusConfiguration.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/WorkflowStatusConfiguration.java index 87683b18c..c08e79ffd 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/WorkflowStatusConfiguration.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/action/WorkflowStatusConfiguration.java @@ -2,20 +2,23 @@ package ru.spcex.clearing.session.state.action; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import ru.spcex.clearing.notification.NotificationSender; import ru.spcex.platform.enumeration.WorkflowStatus; import ru.spcex.platform.imdg.api.ImdgProvider; @Configuration public class WorkflowStatusConfiguration { private final ImdgProvider imdgProvider; + private final NotificationSender notificationSender; - public WorkflowStatusConfiguration(ImdgProvider imdgProvider) { + public WorkflowStatusConfiguration(ImdgProvider imdgProvider, NotificationSender notificationSender) { this.imdgProvider = imdgProvider; + this.notificationSender = notificationSender; } @Bean("pauseWorkflowStatusAction") public UpdateWorkflowStatusAction pauseWorkflowStatusAction() { - return new UpdateWorkflowStatusAction(imdgProvider, WorkflowStatus.Pause); + return new PauseWorkflowStatusWithNotificationAction(imdgProvider, notificationSender); } @Bean("activeWorkflowStatusAction") diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/model/RgsErr.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/model/RgsErr.java new file mode 100644 index 000000000..ecb83ebd9 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/state/model/RgsErr.java @@ -0,0 +1,29 @@ +package ru.spcex.clearing.session.state.model; + +import ru.clearing.classes.statics.data.registry.Registry; + +public class RgsErr { + private Registry rgs; + private String err; + + public RgsErr(Registry rgs, String err) { + this.rgs = rgs; + this.err = err; + } + + public Registry getRgs() { + return rgs; + } + + public void setRgs(Registry rgs) { + this.rgs = rgs; + } + + public String getErr() { + return err; + } + + public void setErr(String err) { + this.err = err; + } +}