This commit is contained in:
parent
fa8cfe268d
commit
936e4613d8
2 changed files with 152 additions and 38 deletions
|
|
@ -3,12 +3,17 @@ package ru.spcex.clearing.session.stage;
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
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.annotation.Qualifier;
|
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.misc.Session;
|
import ru.clearing.classes.statics.data.misc.Session;
|
||||||
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.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.session.stage.impl.*;
|
import ru.spcex.clearing.session.stage.impl.*;
|
||||||
|
import ru.spcex.clearing.session.stage.task.*;
|
||||||
import ru.spcex.platform.enumeration.Section;
|
import ru.spcex.platform.enumeration.Section;
|
||||||
import ru.spcex.platform.enumeration.SessionStatus;
|
import ru.spcex.platform.enumeration.SessionStatus;
|
||||||
import ru.spcex.platform.enumeration.SessionType;
|
import ru.spcex.platform.enumeration.SessionType;
|
||||||
|
|
@ -16,10 +21,11 @@ import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
||||||
import java.util.concurrent.atomic.AtomicBoolean;
|
import java.util.List;
|
||||||
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class PrimaryAuctionBnSession extends QueueConsumer {
|
public class PrimaryAuctionBnSession extends QueueConsumer implements InitializingBean {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
private final BalanceRevise balanceRevise;
|
private final BalanceRevise balanceRevise;
|
||||||
private final DealsPrepare dealsPrepare;
|
private final DealsPrepare dealsPrepare;
|
||||||
|
|
@ -30,13 +36,13 @@ public class PrimaryAuctionBnSession extends QueueConsumer {
|
||||||
private final FormingRegistersOnOS formingRegistersOnOS;
|
private final FormingRegistersOnOS formingRegistersOnOS;
|
||||||
private final FormingPaymentInstruction formingPaymentInstruction;
|
private final FormingPaymentInstruction formingPaymentInstruction;
|
||||||
private final UnlockResources unlockResources;
|
private final UnlockResources unlockResources;
|
||||||
|
private final FinishingSession finishingSession;
|
||||||
private final EndStageNotification endStageNotification;
|
private final EndStageNotification endStageNotification;
|
||||||
|
|
||||||
//
|
|
||||||
private final AtomicBoolean running;
|
|
||||||
private Long sessionId;
|
|
||||||
private Imdg<Session> sessionImdg;
|
private Imdg<Session> sessionImdg;
|
||||||
private final IMessageResolver messageResolver;
|
private final IMessageResolver messageResolver;
|
||||||
|
private final AtomicReference<TaskType> currStage = new AtomicReference<>();
|
||||||
|
private Session currSession;
|
||||||
|
|
||||||
public PrimaryAuctionBnSession(
|
public PrimaryAuctionBnSession(
|
||||||
@Qualifier("createConsumer") Consumer<String, Object> kafkaQueue,
|
@Qualifier("createConsumer") Consumer<String, Object> kafkaQueue,
|
||||||
|
|
@ -49,7 +55,7 @@ public class PrimaryAuctionBnSession extends QueueConsumer {
|
||||||
FormingRegistersOnOS formingRegistersOnOS,
|
FormingRegistersOnOS formingRegistersOnOS,
|
||||||
FormingPaymentInstruction formingPaymentInstruction,
|
FormingPaymentInstruction formingPaymentInstruction,
|
||||||
UnlockResources unlockResources,
|
UnlockResources unlockResources,
|
||||||
EndStageNotification endStageNotification, IMessageResolver messageResolver) {
|
FinishingSession finishingSession, EndStageNotification endStageNotification, IMessageResolver messageResolver) {
|
||||||
super(kafkaQueue);
|
super(kafkaQueue);
|
||||||
this.balanceRevise = balanceRevise;
|
this.balanceRevise = balanceRevise;
|
||||||
this.dealsPrepare = dealsPrepare;
|
this.dealsPrepare = dealsPrepare;
|
||||||
|
|
@ -59,43 +65,151 @@ public class PrimaryAuctionBnSession extends QueueConsumer {
|
||||||
this.formingRegistersOnOS = formingRegistersOnOS;
|
this.formingRegistersOnOS = formingRegistersOnOS;
|
||||||
this.formingPaymentInstruction = formingPaymentInstruction;
|
this.formingPaymentInstruction = formingPaymentInstruction;
|
||||||
this.unlockResources = unlockResources;
|
this.unlockResources = unlockResources;
|
||||||
|
this.finishingSession = finishingSession;
|
||||||
this.endStageNotification = endStageNotification;
|
this.endStageNotification = endStageNotification;
|
||||||
//
|
|
||||||
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
|
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
|
||||||
this.messageResolver = messageResolver;
|
this.messageResolver = messageResolver;
|
||||||
this.running = new AtomicBoolean(false);
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() throws Exception {
|
||||||
|
callback(Object.class)
|
||||||
|
.setConsumer(this::continueSession)
|
||||||
|
.forDestination(Consts.SDF57_PROCESS, callbacks::put);
|
||||||
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
public void runSession() {
|
public void runSession() {
|
||||||
synchronized (this.running) {
|
if (!startSession()) {
|
||||||
if (this.running.get()) {
|
return;
|
||||||
log.info("already running session.id={}", this.sessionId);
|
}
|
||||||
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, companyIdPayload, 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 <T, R> StageResult<R> runStage(TaskType type, ISessionStage stage) {
|
||||||
|
return runStage(type, null, stage);
|
||||||
|
}
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
|
private <T, R> StageResult<R> runStage(TaskType type, T payload, ISessionStage stage) {
|
||||||
|
continueRunning(type);
|
||||||
|
log.info("session.id={} step {} started", currSession.getId(), currStage.get());
|
||||||
|
Task<T> t = new Task<>(type, payload);
|
||||||
|
StageResult<?> stgRes = stage.submit(t);
|
||||||
|
log.info("session.id={} step {} result: {} ",
|
||||||
|
currSession.getId(),
|
||||||
|
currStage.get(),
|
||||||
|
stgRes.success ? "success" : messageResolver.resolve(stgRes.error));
|
||||||
|
if (!stgRes.success) {
|
||||||
|
endSession();
|
||||||
|
throw new StageException();
|
||||||
|
}
|
||||||
|
return (StageResult<R>) stgRes;
|
||||||
|
}
|
||||||
|
|
||||||
|
private boolean startSession() {
|
||||||
|
synchronized (this.currStage) {
|
||||||
|
if (this.currStage.get() != null) {
|
||||||
|
log.info("already running session.id={}", this.currSession.getId());
|
||||||
|
return false;
|
||||||
} else {
|
} else {
|
||||||
this.running.set(true);
|
Session newSession = new Session();
|
||||||
|
newSession.setSection(Section.FOND.getKey());
|
||||||
|
newSession.setSessionType(SessionType.IPOB.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;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Session newSession = new Session();
|
|
||||||
newSession.setSection(Section.FOND.getKey());
|
|
||||||
newSession.setSessionType(SessionType.IPOB.getKey());
|
|
||||||
newSession.setSessionStatus(SessionStatus.CLRN.getKey());
|
|
||||||
//todo companyId/securityId/userId передается из сообщения очереди
|
|
||||||
sessionImdg.insert(newSession);
|
|
||||||
sessionId = newSession.getId();
|
|
||||||
//java.util.function.Consumer<StageResult<?>> logError = stageResult -> {
|
|
||||||
// if (!stageResult.success) {
|
|
||||||
// log.error("stageResult.error={}", messageResolver.resolve(stageResult.error));
|
|
||||||
// newSession.setSessionStatus(SessionStatus.CLOS.getKey());
|
|
||||||
// }
|
|
||||||
//};
|
|
||||||
//{
|
|
||||||
// StageResult<?> reviseResult = balanceRevise.submit(task(TaskType.StartRevise));
|
|
||||||
// logError.accept(reviseResult);
|
|
||||||
//}
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private Task<?> task(TaskType taskType) {
|
private void endSession() {
|
||||||
return new Task<>(taskType, null);
|
synchronized (this.currStage) {
|
||||||
|
this.currStage.set(null);
|
||||||
|
this.currSession = null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void continueRunning(TaskType t) {
|
||||||
|
synchronized (this.currStage) {
|
||||||
|
this.currStage.set(t);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private boolean checkStage(TaskType t) {
|
||||||
|
synchronized (this.currStage) {
|
||||||
|
return this.currStage.get().equals(t);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private static class StageException extends RuntimeException {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
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.ExecutionDeposit;
|
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
||||||
import ru.clearing.classes.statics.data.execution.ExecutionFond;
|
import ru.clearing.classes.statics.data.execution.ExecutionFond;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
|
@ -11,7 +12,6 @@ import ru.spcex.clearing.session.stage.ISessionStage;
|
||||||
import ru.spcex.clearing.session.stage.StageResult;
|
import ru.spcex.clearing.session.stage.StageResult;
|
||||||
import ru.spcex.clearing.session.stage.Task;
|
import ru.spcex.clearing.session.stage.Task;
|
||||||
import ru.spcex.clearing.session.stage.task.DealsPreparePayload;
|
import ru.spcex.clearing.session.stage.task.DealsPreparePayload;
|
||||||
import ru.spcex.platform.classes.base.interfaces.IExecution;
|
|
||||||
import ru.spcex.platform.classes.base.interfaces.WithExchangeExecutionId;
|
import ru.spcex.platform.classes.base.interfaces.WithExchangeExecutionId;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
@ -74,11 +74,11 @@ public class DealsPrepare implements ISessionStage {
|
||||||
ImdgPredicate excFondPrct = prdComposer.apply(execFondPredicates, executionFondImdg);
|
ImdgPredicate excFondPrct = prdComposer.apply(execFondPredicates, executionFondImdg);
|
||||||
Collection<ExecutionDeposit> excDpsts = executionDepositImdg.getCollectionObjectsByPredicate(excDepPrct);
|
Collection<ExecutionDeposit> excDpsts = executionDepositImdg.getCollectionObjectsByPredicate(excDepPrct);
|
||||||
Collection<ExecutionFond> excFonds = executionFondImdg.getCollectionObjectsByPredicate(excFondPrct);
|
Collection<ExecutionFond> excFonds = executionFondImdg.getCollectionObjectsByPredicate(excFondPrct);
|
||||||
List<IExecution> excs = Stream.concat(excDpsts.stream().map(execToInterface()),
|
List<ExecutionCommon> excs = Stream.concat(excDpsts.stream().map(execToInterface()),
|
||||||
excFonds.stream().map(execToInterface()))
|
excFonds.stream().map(execToInterface()))
|
||||||
.sorted(Comparator.comparing(WithExchangeExecutionId::getExchangeExecutionId))
|
.sorted(Comparator.comparing(WithExchangeExecutionId::getExchangeExecutionId))
|
||||||
.toList();
|
.toList();
|
||||||
for (IExecution exc : excs) {
|
for (ExecutionCommon exc : excs) {
|
||||||
exc.setSessionId(sessionId);
|
exc.setSessionId(sessionId);
|
||||||
if (exc instanceof ExecutionDeposit) {
|
if (exc instanceof ExecutionDeposit) {
|
||||||
executionDepositImdg.update((ExecutionDeposit) exc);
|
executionDepositImdg.update((ExecutionDeposit) exc);
|
||||||
|
|
@ -86,7 +86,7 @@ public class DealsPrepare implements ISessionStage {
|
||||||
executionFondImdg.update((ExecutionFond) exc);
|
executionFondImdg.update((ExecutionFond) exc);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
StageResult<List<IExecution>> res = new StageResult<>(null, true);
|
StageResult<List<ExecutionCommon>> res = new StageResult<>(null, true);
|
||||||
res.setStageResult(excs);
|
res.setStageResult(excs);
|
||||||
return res;
|
return res;
|
||||||
}
|
}
|
||||||
|
|
@ -103,7 +103,7 @@ public class DealsPrepare implements ISessionStage {
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
private static <E extends IExecution> Function<E, IExecution> execToInterface() {
|
private static <E extends ExecutionCommon> Function<E, ExecutionCommon> execToInterface() {
|
||||||
return (e) -> e;
|
return (e) -> e;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue