final MKR session duplicate intermediate MKR session
This commit is contained in:
parent
e58c6fa899
commit
29b3f3b188
5 changed files with 238 additions and 4 deletions
|
|
@ -28,6 +28,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
||||||
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
||||||
private final IntermediateMkrSession intermediateMkrSession;
|
private final IntermediateMkrSession intermediateMkrSession;
|
||||||
|
private final FinalMkrSession finalMkrSession;
|
||||||
private final SessionManager sessionManager;
|
private final SessionManager sessionManager;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
|
|
@ -36,7 +37,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
RegistryService registryService,
|
RegistryService registryService,
|
||||||
PrimaryAuctionBnSession primaryAuctionBnSession,
|
PrimaryAuctionBnSession primaryAuctionBnSession,
|
||||||
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
||||||
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, SessionManager sessionManager) {
|
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, SessionManager sessionManager) {
|
||||||
super(kafkaQueue);
|
super(kafkaQueue);
|
||||||
this.clearingService = clearingService;
|
this.clearingService = clearingService;
|
||||||
this.registryService = registryService;
|
this.registryService = registryService;
|
||||||
|
|
@ -45,6 +46,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
||||||
this.primaryAuctionT0Session = primaryAuctionT0Session;
|
this.primaryAuctionT0Session = primaryAuctionT0Session;
|
||||||
this.intermediateMkrSession = intermediateMkrSession;
|
this.intermediateMkrSession = intermediateMkrSession;
|
||||||
|
this.finalMkrSession = finalMkrSession;
|
||||||
this.sessionManager = sessionManager;
|
this.sessionManager = sessionManager;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -77,6 +79,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
secondaryAuctionT0Session.continueSession(req);
|
secondaryAuctionT0Session.continueSession(req);
|
||||||
primaryAuctionB0Session.continueSession(req);
|
primaryAuctionB0Session.continueSession(req);
|
||||||
intermediateMkrSession.continueSession(req);
|
intermediateMkrSession.continueSession(req);
|
||||||
|
finalMkrSession.continueSession(req);
|
||||||
})
|
})
|
||||||
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
|
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,228 @@
|
||||||
|
package ru.spcex.clearing.session.stage;
|
||||||
|
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
|
||||||
|
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
||||||
|
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.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.session.stage.impl.*;
|
||||||
|
import ru.spcex.clearing.session.stage.task.*;
|
||||||
|
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||||
|
import ru.spcex.platform.enumeration.Section;
|
||||||
|
import ru.spcex.platform.enumeration.SessionStatus;
|
||||||
|
import ru.spcex.platform.enumeration.SessionType;
|
||||||
|
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.IMessageResolver;
|
||||||
|
|
||||||
|
import java.time.LocalDate;
|
||||||
|
import java.util.Collection;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.function.Function;
|
||||||
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class FinalMkrSession extends AbstractSession implements InitializingBean {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final BalanceRevise balanceRevise;
|
||||||
|
private final DealsPrepare dealsPrepare;
|
||||||
|
private final RequirementsAndObligationCreation requirementsAndObligationCreation;
|
||||||
|
private final ObligationAdmission obligationsAdmission;
|
||||||
|
|
||||||
|
private final InclusionObligations inclusionObligations;
|
||||||
|
private final InspectionObligations inspectionObligations;
|
||||||
|
private final FormingRegistersOnOS formingRegistersOnOS;
|
||||||
|
private final FormingPaymentInstruction formingPaymentInstruction;
|
||||||
|
private final UnlockResources unlockResources;
|
||||||
|
private final FinishingSession finishingSession;
|
||||||
|
private final EndStageNotification endStageNotification;
|
||||||
|
|
||||||
|
private final Imdg<ExecutionDeposit> executionDepositImdg;
|
||||||
|
private final Supplier<List<String>> marketCodes;
|
||||||
|
|
||||||
|
|
||||||
|
public FinalMkrSession(
|
||||||
|
ImdgProvider imdgProvider,
|
||||||
|
BalanceRevise balanceRevise,
|
||||||
|
DealsPrepare dealsPrepare,
|
||||||
|
RequirementsAndObligationCreation requirementsAndObligationCreation,
|
||||||
|
ObligationAdmission obligationsAdmission,
|
||||||
|
InclusionObligations inclusionObligations,
|
||||||
|
FormingRegistersOnOS formingRegistersOnOS,
|
||||||
|
FormingPaymentInstruction formingPaymentInstruction,
|
||||||
|
UnlockResources unlockResources,
|
||||||
|
FinishingSession finishingSession,
|
||||||
|
EndStageNotification endStageNotification,
|
||||||
|
IMessageResolver messageResolver,
|
||||||
|
InspectionObligations inspectionObligations,
|
||||||
|
@Qualifier("marketCodesForBn") Supplier<List<String>> marketCodes) {
|
||||||
|
super(imdgProvider, messageResolver);
|
||||||
|
this.balanceRevise = balanceRevise;
|
||||||
|
this.dealsPrepare = dealsPrepare;
|
||||||
|
this.requirementsAndObligationCreation = requirementsAndObligationCreation;
|
||||||
|
this.obligationsAdmission = obligationsAdmission;
|
||||||
|
this.inclusionObligations = inclusionObligations;
|
||||||
|
this.formingRegistersOnOS = formingRegistersOnOS;
|
||||||
|
this.formingPaymentInstruction = formingPaymentInstruction;
|
||||||
|
this.unlockResources = unlockResources;
|
||||||
|
this.finishingSession = finishingSession;
|
||||||
|
this.endStageNotification = endStageNotification;
|
||||||
|
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
|
||||||
|
this.inspectionObligations = inspectionObligations;
|
||||||
|
this.marketCodes = marketCodes;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() throws Exception {
|
||||||
|
dealsPrepare.searchForExecutions(ExecutionType.ExecutionDeposit);
|
||||||
|
ImdgPredicateBuilder execFondPb = executionDepositImdg.predicateBuilder();
|
||||||
|
dealsPrepare.addExecutionDepositCondition(execFondPb.regex("firstLegSettlementCode", "(^T0.*$)|(B[^0]\\d*$)"));
|
||||||
|
// dealsPrepare.addExecutionDepositCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0])));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void runSession(BaseRequest<?> req) {
|
||||||
|
if (!startSession()) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
StageResult<?> submit = balanceRevise.submit(new Task<>(TaskType.StartRevise, null));
|
||||||
|
if (!submit.success) {
|
||||||
|
log.error("stage {} error={}", balanceRevise.getClass().getSimpleName(), messageResolver.resolve(submit.error));
|
||||||
|
endSession();
|
||||||
|
} else {
|
||||||
|
log.info("stage BalanceRevise success, waiting for a response from kafka");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public void continueSession(BaseRequest<?> req) {
|
||||||
|
try {
|
||||||
|
if (!isRunning()) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (checkStage(TaskType.StartRevise)) {
|
||||||
|
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
|
||||||
|
StageResult<List<ExecutionCommon>> dealsPreparationResult;
|
||||||
|
{
|
||||||
|
DealsPreparePayload payload = new DealsPreparePayload();
|
||||||
|
payload.setSessionId(currSession.getId());
|
||||||
|
dealsPreparationResult = runStage(TaskType.DealsPrepare, payload, dealsPrepare);
|
||||||
|
//Сначала обрабатываются записи, у которых settlementCode = значению, у которого первый символ ="T", второй =0,
|
||||||
|
// а после settlementCode = значению, у которого первый символ ="B", второй ≠0, последующие =любые цифры
|
||||||
|
dealsPreparationResult.getStageResult().sort((o1, o2) -> {
|
||||||
|
String o1SettlementCode = ((ExecutionDeposit) o1).getFirstLegSettlementCode();
|
||||||
|
String o2SettlementCode = ((ExecutionDeposit) o2).getSecondLegSettlementCode();
|
||||||
|
Function<String, Integer> mapper = (settlementCode) -> {
|
||||||
|
if (settlementCode.startsWith("T")) return -1;
|
||||||
|
else if (settlementCode.startsWith("B")) return 1;
|
||||||
|
else return 0;
|
||||||
|
};
|
||||||
|
return mapper.apply(o1SettlementCode).compareTo(mapper.apply(o2SettlementCode));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
//stage 2
|
||||||
|
runStage(TaskType.RequirementsAndObligationsCreate, dealsPreparationResult.getStageResult(), requirementsAndObligationCreation);
|
||||||
|
//stage 3
|
||||||
|
runStage(TaskType.ObligationsAdmission, currSession.getId(), obligationsAdmission);
|
||||||
|
//stage 4
|
||||||
|
{
|
||||||
|
InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload();
|
||||||
|
inclusionToPoolPayload.setSessionType(currSession.getSessionType());
|
||||||
|
runStage(TaskType.InclusionToPool, inclusionToPoolPayload, inclusionObligations);
|
||||||
|
}
|
||||||
|
//stage 5
|
||||||
|
{
|
||||||
|
InspectionPoolPayload companyIdPayload = new InspectionPoolPayload();
|
||||||
|
companyIdPayload.setProcessedCompanyId(currSession.getCompanyId());
|
||||||
|
runStage(TaskType.InspectionObligations, companyIdPayload, inspectionObligations);
|
||||||
|
}
|
||||||
|
//stage 6
|
||||||
|
runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection<Registry>
|
||||||
|
//stage 7
|
||||||
|
StageResult<Collection<PaymentInstruction>> paymentResult = runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction);
|
||||||
|
if (paymentResult.getStageResult().isEmpty()) {
|
||||||
|
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||||
|
// finishPart(req);
|
||||||
|
}
|
||||||
|
} 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();
|
||||||
|
payload.setSessionId(currSession.getId());
|
||||||
|
runStage(TaskType.FinishingSession, payload, finishingSession);
|
||||||
|
}
|
||||||
|
//stage 11
|
||||||
|
{
|
||||||
|
EndStageNotificationPayload payload = new EndStageNotificationPayload();
|
||||||
|
payload.setSection(currSession.getSection());
|
||||||
|
payload.setSessionId(currSession.getId());
|
||||||
|
runStage(TaskType.EndStageNotification, payload, endStageNotification);
|
||||||
|
}
|
||||||
|
} catch (StageException e) {
|
||||||
|
//already logged
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private boolean startSession() {
|
||||||
|
synchronized (this.currStage) {
|
||||||
|
if (this.currStage.get() != null) {
|
||||||
|
log.info("already running session.id={}", this.currSession.getId());
|
||||||
|
return false;
|
||||||
|
} else {
|
||||||
|
TaskType startStatus = TaskType.StartRevise;
|
||||||
|
Session newSession = new Session();
|
||||||
|
newSession.setSection(section().getKey());
|
||||||
|
newSession.setSessionType(sectionType().getKey());
|
||||||
|
newSession.setSessionStatus(startStatus.getKey());
|
||||||
|
newSession.setWorkflowStatus(SessionStatus.ACTV.getKey());
|
||||||
|
newSession.setClearingDate(LocalDate.now());
|
||||||
|
|
||||||
|
//todo companyId/securityId/userId передается из сообщения очереди
|
||||||
|
sessionImdg.insert(newSession);
|
||||||
|
currSession = newSession;
|
||||||
|
log.info("started new session.id={}", this.currSession.getId());
|
||||||
|
currStage.set(startStatus);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected Section section() {
|
||||||
|
return Section.MKR;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected SessionType sectionType() {
|
||||||
|
return SessionType.FINL;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -84,7 +84,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ
|
||||||
dealsPrepare.searchForExecutions(ExecutionType.ExecutionDeposit);
|
dealsPrepare.searchForExecutions(ExecutionType.ExecutionDeposit);
|
||||||
ImdgPredicateBuilder execFondPb = executionDepositImdg.predicateBuilder();
|
ImdgPredicateBuilder execFondPb = executionDepositImdg.predicateBuilder();
|
||||||
dealsPrepare.addExecutionDepositCondition(execFondPb.regex("firstLegSettlementCode", "(^T0.*$)|(B[^0]\\d*$)"));
|
dealsPrepare.addExecutionDepositCondition(execFondPb.regex("firstLegSettlementCode", "(^T0.*$)|(B[^0]\\d*$)"));
|
||||||
dealsPrepare.addExecutionDepositCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0])));
|
// dealsPrepare.addExecutionDepositCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0])));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
|
||||||
|
|
@ -17,17 +17,19 @@ public class SessionManager {
|
||||||
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
||||||
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
||||||
private final IntermediateMkrSession intermediateMkrSession;
|
private final IntermediateMkrSession intermediateMkrSession;
|
||||||
|
private final FinalMkrSession finalMkrSession;
|
||||||
|
|
||||||
public SessionManager(PrimaryAuctionT0Session primaryAuctionT0Session,
|
public SessionManager(PrimaryAuctionT0Session primaryAuctionT0Session,
|
||||||
PrimaryAuctionBnSession primaryAuctionBnSession,
|
PrimaryAuctionBnSession primaryAuctionBnSession,
|
||||||
PrimaryAuctionB0Session primaryAuctionB0Session,
|
PrimaryAuctionB0Session primaryAuctionB0Session,
|
||||||
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
||||||
IntermediateMkrSession intermediateMkrSession) {
|
IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession) {
|
||||||
this.primaryAuctionT0Session = primaryAuctionT0Session;
|
this.primaryAuctionT0Session = primaryAuctionT0Session;
|
||||||
this.primaryAuctionBnSession = primaryAuctionBnSession;
|
this.primaryAuctionBnSession = primaryAuctionBnSession;
|
||||||
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
||||||
this.secondaryAuctionT0Session = secondaryAuctionT0Session;
|
this.secondaryAuctionT0Session = secondaryAuctionT0Session;
|
||||||
this.intermediateMkrSession = intermediateMkrSession;
|
this.intermediateMkrSession = intermediateMkrSession;
|
||||||
|
this.finalMkrSession = finalMkrSession;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void defineAndStartSession(BaseRequest<LauncherCommandRequest> r) { //task sclr section fond sessiontype trdt
|
public void defineAndStartSession(BaseRequest<LauncherCommandRequest> r) { //task sclr section fond sessiontype trdt
|
||||||
|
|
@ -47,6 +49,7 @@ public class SessionManager {
|
||||||
case IPO0 -> session = primaryAuctionB0Session;
|
case IPO0 -> session = primaryAuctionB0Session;
|
||||||
case TRDT -> session = secondaryAuctionT0Session;
|
case TRDT -> session = secondaryAuctionT0Session;
|
||||||
case MEDM -> session = intermediateMkrSession;
|
case MEDM -> session = intermediateMkrSession;
|
||||||
|
case FINL -> session = finalMkrSession;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (session != null) {
|
if (session != null) {
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
|
||||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
|
||||||
public enum SessionType implements IEnumKey {
|
public enum SessionType implements IEnumKey {
|
||||||
IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), MEDM("MEDM"),
|
IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"), MEDM("MEDM"), FINL("FINL"),
|
||||||
;
|
;
|
||||||
|
|
||||||
SessionType(String key) {
|
SessionType(String key) {
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue