session continue EventReceiver

This commit is contained in:
ialbert 2023-05-30 13:19:27 +03:00
parent 5dfeeaff61
commit 23efed8caa
2 changed files with 62 additions and 22 deletions

View file

@ -13,10 +13,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session; import ru.spcex.clearing.session.stage.*;
import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
import ru.spcex.clearing.session.stage.SecondaryAuctionT0Session;
import ru.spcex.clearing.session.stage.SessionManager;
import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.enumeration.Task;
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
@ -27,6 +24,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
private final ClearingService clearingService; private final ClearingService clearingService;
private final RegistryService registryService; private final RegistryService registryService;
private final PrimaryAuctionBnSession primaryAuctionBnSession; private final PrimaryAuctionBnSession primaryAuctionBnSession;
private final PrimaryAuctionT0Session primaryAuctionT0Session;
private final SecondaryAuctionT0Session secondaryAuctionT0Session; private final SecondaryAuctionT0Session secondaryAuctionT0Session;
private final PrimaryAuctionB0Session primaryAuctionB0Session; private final PrimaryAuctionB0Session primaryAuctionB0Session;
private final SessionManager sessionManager; private final SessionManager sessionManager;
@ -37,13 +35,14 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
RegistryService registryService, RegistryService registryService,
PrimaryAuctionBnSession primaryAuctionBnSession, PrimaryAuctionBnSession primaryAuctionBnSession,
SecondaryAuctionT0Session secondaryAuctionT0Session, SecondaryAuctionT0Session secondaryAuctionT0Session,
PrimaryAuctionB0Session primaryAuctionB0Session, SessionManager sessionManager) { PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, SessionManager sessionManager) {
super(kafkaQueue); super(kafkaQueue);
this.clearingService = clearingService; this.clearingService = clearingService;
this.registryService = registryService; this.registryService = registryService;
this.primaryAuctionBnSession = primaryAuctionBnSession; this.primaryAuctionBnSession = primaryAuctionBnSession;
this.secondaryAuctionT0Session = secondaryAuctionT0Session; this.secondaryAuctionT0Session = secondaryAuctionT0Session;
this.primaryAuctionB0Session = primaryAuctionB0Session; this.primaryAuctionB0Session = primaryAuctionB0Session;
this.primaryAuctionT0Session = primaryAuctionT0Session;
this.sessionManager = sessionManager; this.sessionManager = sessionManager;
} }
@ -72,7 +71,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(Object.class) callback(Object.class)
.setConsumer(req -> { .setConsumer(req -> {
primaryAuctionBnSession.continueSession(req); primaryAuctionBnSession.continueSession(req);
primaryAuctionT0Session.continueSession(req);
secondaryAuctionT0Session.continueSession(req); secondaryAuctionT0Session.continueSession(req);
primaryAuctionB0Session.continueSession(req);
}) })
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put); .forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);

View file

@ -3,10 +3,12 @@ package ru.spcex.clearing.session.stage;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.execution.ExecutionCommon; import ru.clearing.classes.statics.data.execution.ExecutionCommon;
import ru.clearing.classes.statics.data.execution.ExecutionFond; import ru.clearing.classes.statics.data.execution.ExecutionFond;
import ru.clearing.classes.statics.data.misc.Session; import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.session.stage.impl.*; import ru.spcex.clearing.session.stage.impl.*;
@ -21,7 +23,10 @@ import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.enumeration.IMessageResolver;
import java.time.LocalDate;
import java.util.Collection;
import java.util.List; import java.util.List;
import java.util.function.Supplier;
@Service @Service
public class PrimaryAuctionB0Session extends AbstractSession implements InitializingBean { public class PrimaryAuctionB0Session extends AbstractSession implements InitializingBean {
@ -32,6 +37,7 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
private final ObligationAdmission obligationsAdmission; private final ObligationAdmission obligationsAdmission;
private final InclusionObligations inclusionObligations; private final InclusionObligations inclusionObligations;
private final InspectionObligations inspectionObligations;
private final FormingRegistersOnOS formingRegistersOnOS; private final FormingRegistersOnOS formingRegistersOnOS;
private final FormingPaymentInstruction formingPaymentInstruction; private final FormingPaymentInstruction formingPaymentInstruction;
private final UnlockResources unlockResources; private final UnlockResources unlockResources;
@ -39,6 +45,8 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
private final EndStageNotification endStageNotification; private final EndStageNotification endStageNotification;
private final Imdg<ExecutionFond> executionFondImdg; private final Imdg<ExecutionFond> executionFondImdg;
private final Supplier<List<String>> marketCodes;
public PrimaryAuctionB0Session( public PrimaryAuctionB0Session(
ImdgProvider imdgProvider, ImdgProvider imdgProvider,
@ -50,7 +58,11 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
FormingRegistersOnOS formingRegistersOnOS, FormingRegistersOnOS formingRegistersOnOS,
FormingPaymentInstruction formingPaymentInstruction, FormingPaymentInstruction formingPaymentInstruction,
UnlockResources unlockResources, UnlockResources unlockResources,
FinishingSession finishingSession, EndStageNotification endStageNotification, IMessageResolver messageResolver) { FinishingSession finishingSession,
EndStageNotification endStageNotification,
IMessageResolver messageResolver,
InspectionObligations inspectionObligations,
@Qualifier("marketCodesForBn") Supplier<List<String>> marketCodes) {
super(imdgProvider, messageResolver); super(imdgProvider, messageResolver);
this.balanceRevise = balanceRevise; this.balanceRevise = balanceRevise;
this.dealsPrepare = dealsPrepare; this.dealsPrepare = dealsPrepare;
@ -63,6 +75,8 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
this.finishingSession = finishingSession; this.finishingSession = finishingSession;
this.endStageNotification = endStageNotification; this.endStageNotification = endStageNotification;
this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class); this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
this.inspectionObligations = inspectionObligations;
this.marketCodes = marketCodes;
} }
@Override @Override
@ -89,12 +103,23 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
public void continueSession(BaseRequest<?> req) { public void continueSession(BaseRequest<?> req) {
try { try {
if (!checkStage(TaskType.StartRevise)) { if (!isRunning()) {
log.error("cannot continue session, current stage is {}", currStage.get()); return;
throw new StageException();
} }
//stage 0 if (checkStage(TaskType.StartRevise)) {
runStage(TaskType.ContinueRevise, balanceRevise); firstPart(req);
} else if (checkStage(TaskType.FormingPaymentInstruction)) {
finishPart(req);
} else {
log.info("will not continue session, current stage is {}", currStage.get());
}
} catch (StageException e) {
//already logged
}
}
private void firstPart(BaseRequest<?> req) {
try {
//stage 1 //stage 1
StageResult<List<ExecutionCommon>> dealsPreparationResult; StageResult<List<ExecutionCommon>> dealsPreparationResult;
{ {
@ -116,27 +141,38 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
{ {
InspectionPoolPayload companyIdPayload = new InspectionPoolPayload(); InspectionPoolPayload companyIdPayload = new InspectionPoolPayload();
companyIdPayload.setProcessedCompanyId(currSession.getCompanyId()); companyIdPayload.setProcessedCompanyId(currSession.getCompanyId());
runStage(TaskType.InspectionObligations, companyIdPayload, inclusionObligations); runStage(TaskType.InspectionObligations, companyIdPayload, inspectionObligations);
} }
//stage 6 //stage 6
runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection<Registry> runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection<Registry>
//stage 7 //stage 7
runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction); StageResult<Collection<PaymentInstruction>> paymentResult = runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction);
//stage 8 if (paymentResult.getStageResult().isEmpty()) {
{ runStage(TaskType.FormingPaymentInstruction, balanceRevise);
UnlockResourcesPayload unlockResourcesPayload = new UnlockResourcesPayload(); // finishPart(req);
//todo set arguments
runStage(TaskType.UnlockResources, unlockResourcesPayload, unlockResources); //returns Collection<Registry>
} }
//stage 9 } catch (StageException e) {
//already logged
}
}
public void finishPart(BaseRequest<?> req) {
try {
if (!checkStage(TaskType.FormingPaymentInstruction)) {
log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException();
}
//stage 10
{ {
FinishingSessionPayload payload = new FinishingSessionPayload(); FinishingSessionPayload payload = new FinishingSessionPayload();
payload.setSessionId(currSession.getId()); payload.setSessionId(currSession.getId());
runStage(TaskType.FinishingSession, payload, finishingSession); runStage(TaskType.FinishingSession, payload, finishingSession);
} }
//stage 11
{ {
EndStageNotificationPayload payload = new EndStageNotificationPayload(); EndStageNotificationPayload payload = new EndStageNotificationPayload();
payload.setSection(currSession.getSection()); payload.setSection(currSession.getSection());
payload.setSessionId(currSession.getId());
runStage(TaskType.EndStageNotification, payload, endStageNotification); runStage(TaskType.EndStageNotification, payload, endStageNotification);
} }
} catch (StageException e) { } catch (StageException e) {
@ -150,15 +186,18 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
log.info("already running session.id={}", this.currSession.getId()); log.info("already running session.id={}", this.currSession.getId());
return false; return false;
} else { } else {
TaskType startStatus = TaskType.StartRevise;
Session newSession = new Session(); Session newSession = new Session();
newSession.setSection(Section.FOND.getKey()); newSession.setSection(Section.FOND.getKey());
newSession.setSessionType(SessionType.IPO0.getKey()); newSession.setSessionType(SessionType.IPO0.getKey());
newSession.setSessionStatus(SessionStatus.CLRN.getKey()); newSession.setSessionStatus(startStatus.getKey());
newSession.setWorkflowStatus(SessionStatus.ACTV.getKey());
newSession.setClearingDate(LocalDate.now());
//todo companyId/securityId/userId передается из сообщения очереди //todo companyId/securityId/userId передается из сообщения очереди
sessionImdg.insert(newSession); sessionImdg.insert(newSession);
currSession = newSession; currSession = newSession;
log.info("started new session.id={}", this.currSession.getId()); log.info("started new session.id={}", this.currSession.getId());
currStage.set(TaskType.StartRevise); currStage.set(startStatus);
return true; return true;
} }
} }
@ -171,6 +210,6 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
@Override @Override
protected SessionType sectionType() { protected SessionType sectionType() {
return SessionType.IPO0; return SessionType.IPOB;
} }
} }