B0 session
This commit is contained in:
parent
d9cc47aadf
commit
0fe495a1a2
5 changed files with 203 additions and 5 deletions
|
|
@ -9,6 +9,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
|
|||
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.service.QueueConsumer;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
|
||||
import ru.spcex.platform.enumeration.Task;
|
||||
|
||||
|
|
@ -17,15 +18,18 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
private final ClearingService clearingService;
|
||||
private final RegistryService registryService;
|
||||
private final PrimaryAuctionBnSession primaryAuctionBnSession;
|
||||
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
||||
|
||||
public EventsReceiver(Consumer<String, Object> kafkaQueue,
|
||||
ClearingService clearingService,
|
||||
RegistryService registryService,
|
||||
PrimaryAuctionBnSession primaryAuctionBnSession) {
|
||||
PrimaryAuctionBnSession primaryAuctionBnSession,
|
||||
PrimaryAuctionB0Session primaryAuctionB0Session) {
|
||||
super(kafkaQueue);
|
||||
this.clearingService = clearingService;
|
||||
this.registryService = registryService;
|
||||
this.primaryAuctionBnSession = primaryAuctionBnSession;
|
||||
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -42,15 +46,25 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
callback(LauncherCommandRequest.class)
|
||||
.setConsumer(event -> clearingService.executeVerification())
|
||||
.forDestination(Task.getVerification.topic(), callbacks::put);
|
||||
callback(Object.class)
|
||||
.setConsumer(primaryAuctionBnSession::runSession)
|
||||
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
||||
callback(CommonIdRequest.class)
|
||||
.setConsumer(clearingService::continueClearing)
|
||||
.forDestination(Consts.CONTINUE_CLEARING, callbacks::put);
|
||||
|
||||
callback(Object.class)
|
||||
.setConsumer(primaryAuctionB0Session::runSession)
|
||||
.forDestination(Task.startOfB0.topic(), callbacks::put);
|
||||
// callback(Object.class)
|
||||
// .setConsumer(primaryAuctionB0Session::continueSession)
|
||||
// .forDestination(Task.startOfB0.topic(), callbacks::put);
|
||||
|
||||
|
||||
callback(Object.class)
|
||||
.setConsumer(primaryAuctionBnSession::runSession)
|
||||
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
||||
callback(Object.class)
|
||||
.setConsumer(primaryAuctionBnSession::continueSession)
|
||||
.forDestination(Consts.SDF57_PROCESS, callbacks::put);
|
||||
|
||||
callback(Object.class)
|
||||
.setConsumer(event -> clearingService.executeSTrade())
|
||||
.forDestination(Task.getOfTrades.topic(), callbacks::put);
|
||||
|
|
|
|||
|
|
@ -0,0 +1,165 @@
|
|||
package ru.spcex.clearing.session.stage;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.stereotype.Service;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionFond;
|
||||
import ru.clearing.classes.statics.data.misc.Session;
|
||||
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.MarketType;
|
||||
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.util.List;
|
||||
|
||||
@Service
|
||||
public class PrimaryAuctionB0Session 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 FormingRegistersOnOS formingRegistersOnOS;
|
||||
private final FormingPaymentInstruction formingPaymentInstruction;
|
||||
private final UnlockResources unlockResources;
|
||||
private final FinishingSession finishingSession;
|
||||
private final EndStageNotification endStageNotification;
|
||||
|
||||
private final Imdg<ExecutionFond> executionFondImdg;
|
||||
|
||||
public PrimaryAuctionB0Session(
|
||||
ImdgProvider imdgProvider,
|
||||
BalanceRevise balanceRevise,
|
||||
DealsPrepare dealsPrepare,
|
||||
RequirementsAndObligationCreation requirementsAndObligationCreation,
|
||||
ObligationAdmission obligationsAdmission,
|
||||
InclusionObligations inclusionObligations,
|
||||
FormingRegistersOnOS formingRegistersOnOS,
|
||||
FormingPaymentInstruction formingPaymentInstruction,
|
||||
UnlockResources unlockResources,
|
||||
FinishingSession finishingSession, EndStageNotification endStageNotification, IMessageResolver messageResolver) {
|
||||
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.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() throws Exception {
|
||||
dealsPrepare.searchForExecutions(ExecutionType.ExecutionFond);
|
||||
ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder();
|
||||
dealsPrepare.addExecutionFondCondition(execFondPb.regex("settlementCode", "^B0.*$"));
|
||||
dealsPrepare.addExecutionFondCondition(execFondPb.equals("marketType", MarketType.PRMR.getKey()));
|
||||
}
|
||||
|
||||
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 (!checkStage(TaskType.StartRevise)) {
|
||||
log.error("cannot continue session, current stage is {}", currStage.get());
|
||||
throw new StageException();
|
||||
}
|
||||
//stage 0
|
||||
runStage(TaskType.ContinueRevise, balanceRevise);
|
||||
//stage 1
|
||||
StageResult<List<ExecutionCommon>> dealsPreparationResult;
|
||||
{
|
||||
DealsPreparePayload payload = new DealsPreparePayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
dealsPreparationResult = runStage(TaskType.DealsPrepare , payload, dealsPrepare);
|
||||
}
|
||||
//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, inclusionObligations);
|
||||
}
|
||||
//stage 6
|
||||
runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection<Registry>
|
||||
//stage 7
|
||||
runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction);
|
||||
//stage 8
|
||||
{
|
||||
UnlockResourcesPayload unlockResourcesPayload = new UnlockResourcesPayload();
|
||||
//todo set arguments
|
||||
runStage(TaskType.UnlockResources, unlockResourcesPayload, unlockResources); //returns Collection<Registry>
|
||||
}
|
||||
//stage 9
|
||||
{
|
||||
FinishingSessionPayload payload = new FinishingSessionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
runStage(TaskType.FinishingSession, payload, finishingSession);
|
||||
}
|
||||
{
|
||||
EndStageNotificationPayload payload = new EndStageNotificationPayload();
|
||||
payload.setSection(currSession.getSection());
|
||||
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 {
|
||||
Session newSession = new Session();
|
||||
newSession.setSection(Section.FOND.getKey());
|
||||
newSession.setSessionType(SessionType.IPO0.getKey());
|
||||
newSession.setSessionStatus(SessionStatus.CLRN.getKey());
|
||||
//todo companyId/securityId/userId передается из сообщения очереди
|
||||
sessionImdg.insert(newSession);
|
||||
currSession = newSession;
|
||||
log.info("started new session.id={}", this.currSession.getId());
|
||||
currStage.set(TaskType.StartRevise);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,18 @@
|
|||
package ru.spcex.platform.enumeration;
|
||||
|
||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum MarketType implements IEnumKey {
|
||||
PRMR("PRMR");
|
||||
|
||||
private final String key;
|
||||
|
||||
MarketType(String key) {
|
||||
this.key = key;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getKey() {
|
||||
return key;
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
|
|||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum SessionType implements IEnumKey {
|
||||
IPOB("IPOB"),
|
||||
IPOB("IPOB"), IPO0("IPO0"),
|
||||
;
|
||||
|
||||
SessionType(String key) {
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ public enum Task implements IEnumKey {
|
|||
getVerification("GVER"),// Запуск сверки
|
||||
@Deprecated /* todo GBLD удаляется по CLS-267, CLS-275 */ getBalance("GBLD"),// Поступление средств
|
||||
startOfClearing("SCLR"),// Запуск клиринговой сессии
|
||||
startOfB0("IPO0"),// Запуск клиринговой сессии
|
||||
startOfPreClearing("SPRC"),// Запуск преклиринга
|
||||
startPostClearing("SPOC"),// Запуск постклиринга
|
||||
createOrder("CORD"),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue