Obligation Admission PAUSE notification

This commit is contained in:
ialbert 2025-07-11 14:12:28 +03:00
parent c1c7e168ac
commit e47ff42602
6 changed files with 103 additions and 3 deletions

View file

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

View file

@ -11,6 +11,7 @@ public enum DataEnum {
dealsPreparedCurrency,
counterPartyId, //Long
obligationAdmissionStashedRgs, //Map<Long, Registry>
erroneousRegistries, //List<RgsErr>
paymentInstructionReturnMkr, //List<PaymentInstruction>
paymentInstructionSecurity, //List<PaymentInstruction>
paymentsWereCreatedDepositReturn, //Collection<PaymentInstruction>

View file

@ -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<Long, Registry> rgsToUpdate = new HashMap<>();
List<RgsErr> errs = new ArrayList<>();
for (Map.Entry<GroupRgsKey, List<Registry>> grpEntry : byGroups.entrySet()) {
EnumMessage groupError = null;
List<Registry> rgsGroup = grpEntry.getValue();
@ -86,6 +88,7 @@ public class ObligationAdmissionAction extends AbstractSessionActionForOkErrorHa
Optional<EnumMessage> 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);

View file

@ -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<TaskType, SsnEvent> ctx) {
super.actualExecute(ctx);
@SuppressWarnings("unchecked")
List<RgsErr> 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<TaskType, SsnEvent> ctx) {
State<TaskType, SsnEvent> 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";
}
}
}
}

View file

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

View file

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