send notification and return error if session already exist

This commit is contained in:
etreschenkov 2025-07-16 19:43:43 +03:00
parent 627f8caf27
commit 9402218596
2 changed files with 28 additions and 5 deletions

View file

@ -126,7 +126,11 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
callback(LauncherCommandRequest.class)
.setFunction(r -> {
try {
stateMachineWrapper.defineAndStartSession(r);
} catch (ValidationException e) {
return makeErrorResponse(r, e);
}
return null;
})
.forDestination(Task.startOfClearing.topic(), callbacks::put);

View file

@ -12,7 +12,9 @@ import org.springframework.statemachine.StateMachine;
import org.springframework.statemachine.config.StateMachineFactory;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.misc.Session;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.notification.NotificationSender;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.session.stage.TaskType;
@ -20,13 +22,18 @@ import ru.spcex.clearing.session.state.DataEnum;
import ru.spcex.clearing.session.state.SsnEvent;
import ru.spcex.clearing.session.state.interceptor.SessionStatusChangingInterceptor;
import ru.spcex.clearing.session.state.listener.MachineStopListener;
import ru.spcex.platform.enumeration.ObjectType;
import ru.spcex.platform.enumeration.Priority;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.SessionType;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.error.ValidationException;
@Component
public class SessionStateMachineWrapper {
@ -36,12 +43,16 @@ public class SessionStateMachineWrapper {
private final SessionStatusChangingInterceptor statusChangingInterceptor;
private StateMachine<TaskType, SsnEvent> currentSession;
private final Imdg<Session> ssnImdg;
private final NotificationSender notification;
private final IMessageResolver msgs;
@Autowired
public SessionStateMachineWrapper(
ImdgProvider imdgProvider,
Map<String, StateMachineFactory<TaskType, SsnEvent>> factories,
SessionStatusChangingInterceptor statusChangingInterceptor
SessionStatusChangingInterceptor statusChangingInterceptor,
NotificationSender notification,
IMessageResolver msgs
) {
this.factories = factories
.entrySet()
@ -54,9 +65,11 @@ public class SessionStateMachineWrapper {
.collect(Collectors.toMap(AbstractMap.SimpleEntry::getKey, AbstractMap.SimpleEntry::getValue));
this.statusChangingInterceptor = statusChangingInterceptor;
this.ssnImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
this.notification = notification;
this.msgs = msgs;
}
public synchronized void defineAndStartSession(BaseRequest<LauncherCommandRequest> r) {
public synchronized void defineAndStartSession(BaseRequest<LauncherCommandRequest> r) throws ValidationException {
LauncherCommandRequest payload = r.getRequestPayload();
SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, payload.getSessionType());
Section section = IEnumKey.getEnumByKey(Section.class, payload.getSection());
@ -69,6 +82,7 @@ public class SessionStateMachineWrapper {
log.info("no spring state machine factory for {}", sessionType);
return;
}
EnumMessage err = null;
if (this.currentSession != null) {
Session session = currentSession
.getExtendedState()
@ -76,15 +90,20 @@ public class SessionStateMachineWrapper {
Long id = session != null ? session.getId() : null;
String ssnType = session != null ? session.getSessionType() : null;
log.warn("Session already running: type={} id={}", ssnType, id);
return;
err = new EnumMessage(ClearingError.ActiveSessionIsPresent, String.valueOf(id));
}
Session activeSession = getActiveSession().orElse(null);
if (activeSession != null) {
log.warn("Session already running: type={} id={}",
activeSession.getSessionType(),
activeSession.getId());
return;
err = new EnumMessage(ClearingError.ActiveSessionIsPresent, String.valueOf(activeSession.getId()));
}
if (err != null) {
notification.sendNotification(ObjectType.session, msgs.resolve(err), Priority.HIGH);
throw new ValidationException(err);
}
this.currentSession = build(factory);
this.currentSession.addStateListener(new MachineStopListener(
this::clearCurrentMachine