FinalMkrSession state machine integration
This commit is contained in:
parent
68496730c6
commit
ea77a68fb6
7 changed files with 54 additions and 13 deletions
|
|
@ -99,6 +99,7 @@ public class FinalSessionStateMachineConfig
|
||||||
this.endStageNotificationAction = endStageNotificationAction;
|
this.endStageNotificationAction = endStageNotificationAction;
|
||||||
this.finishingSessionAction.setSessionType(SESSION_TYPE);
|
this.finishingSessionAction.setSessionType(SESSION_TYPE);
|
||||||
this.finishingSessionAction.setSection(SECTION);
|
this.finishingSessionAction.setSection(SECTION);
|
||||||
|
this.finishingSessionAction.setPr("3");
|
||||||
|
|
||||||
//stages settings:
|
//stages settings:
|
||||||
ImdgPredicateBuilder rgsPrctBuilder = rgsImdg.predicateBuilder();
|
ImdgPredicateBuilder rgsPrctBuilder = rgsImdg.predicateBuilder();
|
||||||
|
|
@ -149,7 +150,7 @@ public class FinalSessionStateMachineConfig
|
||||||
public void configure(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
public void configure(StateMachineTransitionConfigurer<TaskType, SsnEvent> transitions) throws Exception {
|
||||||
transitions
|
transitions
|
||||||
.withExternal()
|
.withExternal()
|
||||||
.event(SsnEvent.Sdf57Processed)
|
.event(SsnEvent.SDF_57)
|
||||||
.source(TaskType.StartRevise).target(TaskType.StartRevisePart1)
|
.source(TaskType.StartRevise).target(TaskType.StartRevisePart1)
|
||||||
.action(reviseStage1Action)
|
.action(reviseStage1Action)
|
||||||
.and()
|
.and()
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,8 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service;
|
||||||
|
|
||||||
import java.time.LocalTime;
|
import java.time.LocalTime;
|
||||||
|
import java.util.Collection;
|
||||||
|
import java.util.stream.Stream;
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
import org.apache.kafka.clients.producer.Producer;
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
|
|
@ -51,6 +53,8 @@ import ru.spcex.clearing.session.stage.SessionTerminator;
|
||||||
import ru.spcex.clearing.session.stage.TaskType;
|
import ru.spcex.clearing.session.stage.TaskType;
|
||||||
import ru.spcex.clearing.session.stage.UnitedSession;
|
import ru.spcex.clearing.session.stage.UnitedSession;
|
||||||
import ru.spcex.clearing.session.stage.impl.BalanceRevise;
|
import ru.spcex.clearing.session.stage.impl.BalanceRevise;
|
||||||
|
import ru.spcex.clearing.session.state.SsnEvent;
|
||||||
|
import ru.spcex.clearing.session.state.factory.SessionStateMachineWrapper;
|
||||||
import ru.spcex.clearing.statement.StatementServiceV2;
|
import ru.spcex.clearing.statement.StatementServiceV2;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
@ -74,6 +78,7 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
|
||||||
private final UnitedSession unitedSession;
|
private final UnitedSession unitedSession;
|
||||||
private final ReturnDepositSession returnDepositSession;
|
private final ReturnDepositSession returnDepositSession;
|
||||||
private final SessionManager sessionManager;
|
private final SessionManager sessionManager;
|
||||||
|
private final SessionStateMachineWrapper stateMachineWrapper;
|
||||||
private final Sdf06Executor sdf06Executor;
|
private final Sdf06Executor sdf06Executor;
|
||||||
private final Sdf10Executor sdf10Executor;
|
private final Sdf10Executor sdf10Executor;
|
||||||
private final BalanceRevise balanceRevise;
|
private final BalanceRevise balanceRevise;
|
||||||
|
|
@ -93,7 +98,7 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
|
||||||
PrimaryAuctionBnSession primaryAuctionBnSession,
|
PrimaryAuctionBnSession primaryAuctionBnSession,
|
||||||
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
||||||
CurrencySession currencySession,
|
CurrencySession currencySession,
|
||||||
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, PaymSession paymSession, UnitedSession unitedSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
|
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, UnitedSession unitedSession, ReturnDepositSession returnDepositSession, PaymSession paymSession, SessionManager sessionManager, SessionStateMachineWrapper stateMachineWrapper,
|
||||||
Sdf06Executor sdf06Executor,
|
Sdf06Executor sdf06Executor,
|
||||||
Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, SessionTerminator sessionTerminator, PaymentInstructionOutboundService pmtOutboundService, ExecutionCurrencyComponent execCurrUpload, ExecutionDepositComponent execDepUpload, ExecutionFondComponent execFondUpload) {
|
Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, SessionTerminator sessionTerminator, PaymentInstructionOutboundService pmtOutboundService, ExecutionCurrencyComponent execCurrUpload, ExecutionDepositComponent execDepUpload, ExecutionFondComponent execFondUpload) {
|
||||||
super(kafkaQueue, kafkaResponseQueue);
|
super(kafkaQueue, kafkaResponseQueue);
|
||||||
|
|
@ -112,6 +117,7 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
|
||||||
this.unitedSession = unitedSession;
|
this.unitedSession = unitedSession;
|
||||||
this.returnDepositSession = returnDepositSession;
|
this.returnDepositSession = returnDepositSession;
|
||||||
this.sessionManager = sessionManager;
|
this.sessionManager = sessionManager;
|
||||||
|
this.stateMachineWrapper = stateMachineWrapper;
|
||||||
this.sdf06Executor = sdf06Executor;
|
this.sdf06Executor = sdf06Executor;
|
||||||
this.sdf10Executor = sdf10Executor;
|
this.sdf10Executor = sdf10Executor;
|
||||||
this.balanceRevise = balanceRevise;
|
this.balanceRevise = balanceRevise;
|
||||||
|
|
@ -149,6 +155,7 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
|
||||||
.setFunction(r -> {
|
.setFunction(r -> {
|
||||||
try {
|
try {
|
||||||
sessionManager.defineAndStartSession(r);
|
sessionManager.defineAndStartSession(r);
|
||||||
|
stateMachineWrapper.defineAndStartSession(r);
|
||||||
} catch (ValidationException ve) {
|
} catch (ValidationException ve) {
|
||||||
return makeErrorResponse(r, ve);
|
return makeErrorResponse(r, ve);
|
||||||
}
|
}
|
||||||
|
|
@ -170,10 +177,23 @@ public class EventsReceiver extends QueueConsumerV2 implements InitializingBean
|
||||||
currencySession.continueSession(req);
|
currencySession.continueSession(req);
|
||||||
primaryAuctionB0Session.continueSession(req);
|
primaryAuctionB0Session.continueSession(req);
|
||||||
intermediateMkrSession.continueSession(req);
|
intermediateMkrSession.continueSession(req);
|
||||||
finalMkrSession.continueSession(req);
|
//finalMkrSession.continueSession(req);
|
||||||
unitedSession.continueSession(req);
|
unitedSession.continueSession(req);
|
||||||
returnDepositSession.continueSession(req);
|
returnDepositSession.continueSession(req);
|
||||||
paymSession.continueSession(req);
|
paymSession.continueSession(req);
|
||||||
|
SessionContinueEvent payload = req.getRequestPayload();
|
||||||
|
Stream
|
||||||
|
.ofNullable(payload.getSdfType())
|
||||||
|
.flatMap(Collection::stream)
|
||||||
|
.forEach(sdf -> {
|
||||||
|
try {
|
||||||
|
SsnEvent ssnEvent = SsnEvent.valueOf(sdf.name());
|
||||||
|
stateMachineWrapper.sendEvent(ssnEvent);
|
||||||
|
} catch (IllegalArgumentException e) {
|
||||||
|
log.info("sdf {} event skipped for state machine", sdf);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
;
|
||||||
})
|
})
|
||||||
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
|
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -77,6 +77,10 @@ public class SessionManager {
|
||||||
log.warn("Unknown sessionType: {} or section: {}", section, sessionType);
|
log.warn("Unknown sessionType: {} or section: {}", section, sessionType);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
if (SessionType.FINL.equals(sessionType)) {
|
||||||
|
log.info("falling back to spring state machine");
|
||||||
|
return;
|
||||||
|
}
|
||||||
BaseRequest<?> baseRequest = new BaseRequest<>();
|
BaseRequest<?> baseRequest = new BaseRequest<>();
|
||||||
AbstractSession session = sessionByType(sessionType);
|
AbstractSession session = sessionByType(sessionType);
|
||||||
if (session != null) {
|
if (session != null) {
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
package ru.spcex.clearing.session.state;
|
package ru.spcex.clearing.session.state;
|
||||||
|
|
||||||
public enum ActionHeader {
|
public enum ActionHeader {
|
||||||
sdf56FromTime, pr;
|
sdf56FromTime;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,6 @@
|
||||||
package ru.spcex.clearing.session.state;
|
package ru.spcex.clearing.session.state;
|
||||||
|
|
||||||
public enum SsnEvent {
|
public enum SsnEvent {
|
||||||
Revise,
|
|
||||||
Sdf57Processed,
|
|
||||||
SdfReceived,
|
|
||||||
SDF_01,
|
SDF_01,
|
||||||
SDF_57,
|
SDF_57,
|
||||||
SDF_04,
|
SDF_04,
|
||||||
|
|
|
||||||
|
|
@ -35,7 +35,6 @@ import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.clearing.service.Sdf05Sender;
|
import ru.spcex.clearing.service.Sdf05Sender;
|
||||||
import ru.spcex.clearing.service.Sdf14Sender;
|
import ru.spcex.clearing.service.Sdf14Sender;
|
||||||
import ru.spcex.clearing.session.stage.TaskType;
|
import ru.spcex.clearing.session.stage.TaskType;
|
||||||
import ru.spcex.clearing.session.state.ActionHeader;
|
|
||||||
import ru.spcex.clearing.session.state.DataEnum;
|
import ru.spcex.clearing.session.state.DataEnum;
|
||||||
import ru.spcex.clearing.session.state.SsnEvent;
|
import ru.spcex.clearing.session.state.SsnEvent;
|
||||||
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||||
|
|
@ -74,6 +73,7 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl
|
||||||
private final Sdf14Sender sdf14Sender;
|
private final Sdf14Sender sdf14Sender;
|
||||||
private SessionType sessionType;
|
private SessionType sessionType;
|
||||||
private Section section;
|
private Section section;
|
||||||
|
private String pr;
|
||||||
|
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
|
|
@ -98,6 +98,11 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl
|
||||||
public void setSection(Section section) {
|
public void setSection(Section section) {
|
||||||
this.section = section;
|
this.section = section;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public void setPr(String pr) {
|
||||||
|
this.pr = pr;
|
||||||
|
}
|
||||||
|
|
||||||
public void setSessionType(SessionType sessionType) {
|
public void setSessionType(SessionType sessionType) {
|
||||||
this.sessionType = sessionType;
|
this.sessionType = sessionType;
|
||||||
}
|
}
|
||||||
|
|
@ -106,7 +111,9 @@ public class FinishingSessionAction extends AbstractSessionActionForOkErrorHandl
|
||||||
@Override
|
@Override
|
||||||
protected void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
|
protected void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
|
||||||
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
|
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
|
||||||
String pr = (String) ctx.getMessageHeader(ActionHeader.pr);
|
if (pr == null) {
|
||||||
|
throw new IllegalStateException("'pr' not set up");
|
||||||
|
}
|
||||||
Instant now = Instant.now();
|
Instant now = Instant.now();
|
||||||
//установка CLRD для обработанных регистров
|
//установка CLRD для обработанных регистров
|
||||||
Collection<Registry> allLiabilitiesAndOKClaims = selectClaimsAndLiabilities(sessionId);
|
Collection<Registry> allLiabilitiesAndOKClaims = selectClaimsAndLiabilities(sessionId);
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,6 @@ package ru.spcex.clearing.session.state.factory;
|
||||||
|
|
||||||
import java.util.AbstractMap;
|
import java.util.AbstractMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.NoSuchElementException;
|
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
@ -57,6 +56,10 @@ public class SessionStateMachineWrapper {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
StateMachineFactory<TaskType, SsnEvent> factory = getByType(sessionType);
|
StateMachineFactory<TaskType, SsnEvent> factory = getByType(sessionType);
|
||||||
|
if (factory == null) {
|
||||||
|
log.info("no spring state machine factory for {}", sessionType);
|
||||||
|
return;
|
||||||
|
}
|
||||||
if (this.currentSession != null) {
|
if (this.currentSession != null) {
|
||||||
Session session = currentSession
|
Session session = currentSession
|
||||||
.getExtendedState()
|
.getExtendedState()
|
||||||
|
|
@ -82,6 +85,15 @@ public class SessionStateMachineWrapper {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public synchronized void sendEvent(SsnEvent event) {
|
||||||
|
if (this.currentSession != null) {
|
||||||
|
log.info("sending an event to session: {}", event);
|
||||||
|
this.currentSession.sendEvent(event);
|
||||||
|
} else {
|
||||||
|
log.warn("there is no currently active session");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private synchronized void clearCurrentMachine() {
|
private synchronized void clearCurrentMachine() {
|
||||||
this.currentSession.stop();
|
this.currentSession.stop();
|
||||||
this.currentSession = null;
|
this.currentSession = null;
|
||||||
|
|
@ -89,9 +101,9 @@ public class SessionStateMachineWrapper {
|
||||||
|
|
||||||
private StateMachineFactory<TaskType, SsnEvent> getByType(SessionType sessionType) {
|
private StateMachineFactory<TaskType, SsnEvent> getByType(SessionType sessionType) {
|
||||||
StateMachineFactory<TaskType, SsnEvent> factory = factories.get(sessionType);
|
StateMachineFactory<TaskType, SsnEvent> factory = factories.get(sessionType);
|
||||||
if (factory == null) {
|
//if (factory == null) {
|
||||||
throw new NoSuchElementException("no state machine for session: " + sessionType);
|
// throw new NoSuchElementException("no state machine for session: " + sessionType);
|
||||||
}
|
//}
|
||||||
return factory;
|
return factory;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue