added session monitor for sdf04/sdf13/sdf01/sdf57
This commit is contained in:
parent
55f66994b3
commit
7caa78edef
19 changed files with 276 additions and 73 deletions
|
|
@ -11,6 +11,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueEvent;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryReturnDepositRequest;
|
||||
|
|
@ -88,7 +89,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
.setConsumer(sessionManager::defineAndStartSession)
|
||||
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
||||
|
||||
callback(Object.class)
|
||||
callback(SessionContinueEvent.class)
|
||||
.setConsumer(req -> {
|
||||
primaryAuctionBnSession.continueSession(req);
|
||||
primaryAuctionT0Session.continueSession(req);
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueEvent;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.service.executors.AbstractExecutor;
|
||||
|
|
@ -94,7 +94,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
processSdf57(pair.getSecond());
|
||||
reviser.doRevise(pair.getFirst().getGroupId());
|
||||
//теперь можем продолжить сессию с шага 1
|
||||
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
|
||||
SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_01, SdfTable.SDF_57);
|
||||
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
|
||||
pairOfSdfRequest.remove(key);
|
||||
log.info("pair sdf01/sdf57 processed successfully");
|
||||
|
|
@ -226,7 +226,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
AbstractExecutor service = executorsMap.get(SdfTable.SDF_04);
|
||||
if (service != null) {
|
||||
Result res = service.execute(sdfGroup, statementRequest);
|
||||
// finishSendCommand(res, service, statementRequest);
|
||||
SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_04);
|
||||
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
|
||||
} else {
|
||||
log.warn("Executor for SDF_04 not set");
|
||||
}
|
||||
|
|
@ -261,6 +262,8 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
|||
"generationId", statementRequest.getGroupId()));
|
||||
AbstractExecutor service = executorsMap.get(SdfTable.SDF_13);
|
||||
Result res = service.execute(sdfGroup, statementRequest);
|
||||
SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_13);
|
||||
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
|
||||
// finishSendCommand(res, service, statementRequest);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -18,9 +18,7 @@ import ru.clearing.classes.statics.data.security.Security;
|
|||
import ru.clearing.classes.statics.data.statement.Statement;
|
||||
import ru.spcex.clearing.error.ClearingError;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionBnRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||
import ru.spcex.clearing.service.AnltSearcher;
|
||||
import ru.spcex.clearing.service.LoggingService;
|
||||
|
|
@ -113,9 +111,6 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
|||
|
||||
@Override
|
||||
public void sendCommand(KafkaSender kafkaSender, Result result) {
|
||||
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
|
||||
continueSessionBn.setGenerationId(result.getGenerationId());
|
||||
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
|
||||
}
|
||||
|
||||
//V - Изменение statement по sDf57
|
||||
|
|
|
|||
|
|
@ -5,6 +5,8 @@ import org.slf4j.LoggerFactory;
|
|||
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.platform.messaging.domain.cud.clearing.SessionContinueEvent;
|
||||
import ru.spcex.clearing.session.stage.monitor.SessionMonitor;
|
||||
import ru.spcex.platform.enumeration.Section;
|
||||
import ru.spcex.platform.enumeration.SessionStatus;
|
||||
import ru.spcex.platform.enumeration.SessionType;
|
||||
|
|
@ -15,8 +17,11 @@ import ru.spcex.platform.utils.enumeration.IEnumKey;
|
|||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.util.Arrays;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
public abstract class AbstractSession {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
|
|
@ -35,6 +40,25 @@ public abstract class AbstractSession {
|
|||
this.messageResolver = messageResolver;
|
||||
}
|
||||
|
||||
protected boolean isMonitorPassed(SessionMonitor m, Object event) {
|
||||
Objects.requireNonNull(m);
|
||||
boolean allSdfsCame = m.isMonitorPassed((SessionContinueEvent) event);
|
||||
if (!allSdfsCame) {
|
||||
log.info("cannot continue session {} {} yet (monitor {})", section(), sessionType(), m.allConditions());
|
||||
} else {
|
||||
log.info("session {} {} monitor {} passed", section(), sessionType(), m.allConditions());
|
||||
}
|
||||
return allSdfsCame;
|
||||
}
|
||||
|
||||
protected String logMonitors(SessionMonitor... monitors) {
|
||||
if (monitors == null || monitors.length == 0) return "<no monitors>";
|
||||
return Arrays.stream(monitors)
|
||||
.filter(Objects::nonNull)
|
||||
.map(SessionMonitor::allConditions)
|
||||
.collect(Collectors.joining("/"));
|
||||
}
|
||||
|
||||
protected <T, R> StageResult<R> runStage(TaskType type, ISessionStage stage) {
|
||||
return runStage(type, null, stage);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,6 +13,8 @@ import ru.clearing.classes.statics.data.registry.Registry;
|
|||
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.monitor.SessionMonitor;
|
||||
import ru.spcex.clearing.session.stage.monitor.SessionMonitorFactory;
|
||||
import ru.spcex.clearing.session.stage.task.*;
|
||||
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||
import ru.spcex.platform.enumeration.RegistryStatus;
|
||||
|
|
@ -51,7 +53,9 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
|
|||
private final Imdg<ExecutionDeposit> executionDepositImdg;
|
||||
private final Supplier<List<String>> marketCodes;
|
||||
private final Imdg<Registry> registryImdg;
|
||||
|
||||
private SessionMonitor firstReviseMonitor;
|
||||
private SessionMonitor afterPaymentsSdf4Monitor;
|
||||
private SessionMonitor afterPaymentsReviseMonitor;
|
||||
|
||||
public FinalMkrSession(
|
||||
ImdgProvider imdgProvider,
|
||||
|
|
@ -120,6 +124,7 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
|
|||
endSession();
|
||||
} else {
|
||||
log.info("stage BalanceRevise success, waiting for a response from kafka");
|
||||
this.firstReviseMonitor = SessionMonitorFactory.waitRevise();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -128,38 +133,38 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
|
|||
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());
|
||||
log.info("session is running, stage {}, monitors: {}",
|
||||
currStage.get(),
|
||||
logMonitors(firstReviseMonitor, afterPaymentsSdf4Monitor, afterPaymentsReviseMonitor));
|
||||
if (firstReviseMonitor != null && isMonitorPassed(firstReviseMonitor, req.getRequestPayload())) {
|
||||
firstReviseMonitor = null;
|
||||
firstPart();
|
||||
}
|
||||
if (afterPaymentsSdf4Monitor != null && isMonitorPassed(afterPaymentsSdf4Monitor, req.getRequestPayload())) {
|
||||
afterPaymentsSdf4Monitor = null;
|
||||
sendSdf56();
|
||||
}
|
||||
if (afterPaymentsReviseMonitor != null && isMonitorPassed(afterPaymentsReviseMonitor, req.getRequestPayload())) {
|
||||
afterPaymentsReviseMonitor = null;
|
||||
finishPart();
|
||||
}
|
||||
} catch (StageException e) {
|
||||
//already logged
|
||||
}
|
||||
}
|
||||
|
||||
private void firstPart(BaseRequest<?> req) {
|
||||
private void firstPart() {
|
||||
try {
|
||||
if (!checkStage(TaskType.StartRevise)) {
|
||||
log.error("cannot continue session, current stage is {}", currStage.get());
|
||||
throw new StageException();
|
||||
}
|
||||
//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);
|
||||
|
|
@ -188,7 +193,7 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
|
|||
}
|
||||
//stage 7
|
||||
|
||||
StageResult<Collection<PaymentInstruction>> returnsPayment = null;
|
||||
StageResult<Collection<PaymentInstruction>> returnsPayment;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
|
|
@ -196,7 +201,7 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
|
|||
returnsPayment = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionReturn);
|
||||
}
|
||||
|
||||
StageResult<Collection<PaymentInstruction>> paymentResult = null;
|
||||
StageResult<PaymentInfo> paymentResult;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
|
|
@ -204,16 +209,24 @@ public class FinalMkrSession extends AbstractSession implements InitializingBean
|
|||
//stage 7
|
||||
paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDeals);
|
||||
}
|
||||
if (paymentResult != null && paymentResult.getStageResult().isEmpty()) {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
// finishPart(req);
|
||||
if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
|
||||
log.info("no payment instructions were created, sending SDF56");
|
||||
sendSdf56();
|
||||
} else {
|
||||
log.info("created {} PaymentInstructions, waiting for SDF04", paymentResult.getStageResult().getPaymentInstructions().size());
|
||||
this.afterPaymentsSdf4Monitor = SessionMonitorFactory.paymentsWereCreated(section());
|
||||
}
|
||||
} catch (StageException e) {
|
||||
//already logged
|
||||
}
|
||||
}
|
||||
|
||||
public void finishPart(BaseRequest<?> req) {
|
||||
private void sendSdf56() {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise();
|
||||
}
|
||||
|
||||
public void finishPart() {
|
||||
try {
|
||||
if (!checkStage(TaskType.FormingPaymentInstruction)) {
|
||||
log.error("cannot continue session, current stage is {}", currStage.get());
|
||||
|
|
|
|||
|
|
@ -184,7 +184,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ
|
|||
returnsPayment = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionReturn);
|
||||
}
|
||||
|
||||
StageResult<Collection<PaymentInstruction>> paymentResult = null;
|
||||
StageResult<PaymentInfo> paymentResult = null;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
|
|
@ -192,7 +192,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ
|
|||
//stage 7
|
||||
paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDeals);
|
||||
}
|
||||
if (paymentResult != null && paymentResult.getStageResult().isEmpty()) {
|
||||
if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
// finishPart(req);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,7 +8,6 @@ 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.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.*;
|
||||
|
|
@ -24,7 +23,6 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
|||
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
|
|
@ -159,7 +157,7 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
|
|||
payload.setSessionId(currSession.getId());
|
||||
runStage(TaskType.FormingRegistersOnOS, payload, formingRegistersOnOS); //returns Collection<Registry>
|
||||
}
|
||||
StageResult<Collection<PaymentInstruction>> paymentResult = null;
|
||||
StageResult<PaymentInfo> paymentResult = null;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
|
|
@ -176,7 +174,7 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
|
|||
//
|
||||
// paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDealsFinalMkr);
|
||||
// }
|
||||
if (paymentResult != null && paymentResult.getStageResult().isEmpty()) {
|
||||
if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
// finishPart(req);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,7 +8,6 @@ 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.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.*;
|
||||
|
|
@ -24,7 +23,6 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
|||
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
|
|
@ -161,7 +159,7 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali
|
|||
runStage(TaskType.FormingRegistersOnOS, payload, formingRegistersOnOS); //returns Collection<Registry>
|
||||
}
|
||||
//stage 7
|
||||
StageResult<Collection<PaymentInstruction>> paymentResult = null;
|
||||
StageResult<PaymentInfo> paymentResult = null;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
|
|
@ -178,7 +176,7 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali
|
|||
//
|
||||
// paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDealsFinalMkr);
|
||||
// }
|
||||
if (paymentResult != null && paymentResult.getStageResult().isEmpty()) {
|
||||
if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
// finishPart(req);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,7 +8,6 @@ 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.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.*;
|
||||
|
|
@ -24,7 +23,6 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
|||
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
|
|
@ -160,7 +158,7 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali
|
|||
payload.setSessionId(currSession.getId());
|
||||
runStage(TaskType.FormingRegistersOnOS, payload, formingRegistersOnOS); //returns Collection<Registry>
|
||||
}
|
||||
StageResult<Collection<PaymentInstruction>> paymentResult = null;
|
||||
StageResult<PaymentInfo> paymentResult = null;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
|
|
@ -177,7 +175,7 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali
|
|||
//
|
||||
// paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDealsFinalMkr);
|
||||
// }
|
||||
if (paymentResult != null && paymentResult.getStageResult().isEmpty()) {
|
||||
if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
// finishPart(req);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,7 +8,6 @@ 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.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.*;
|
||||
|
|
@ -24,7 +23,6 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
|||
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
|
|
@ -157,14 +155,14 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia
|
|||
runStage(TaskType.FormingRegistersOnOS, payload, formingRegistersOnOS); //returns Collection<Registry>
|
||||
}
|
||||
//stage 7
|
||||
StageResult<Collection<PaymentInstruction>> paymentResult = null;
|
||||
StageResult<PaymentInfo> paymentResult = null;
|
||||
{
|
||||
FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
|
||||
payload.setSessionId(currSession.getId());
|
||||
paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets);
|
||||
}
|
||||
|
||||
if (paymentResult != null && paymentResult.getStageResult().isEmpty()) {
|
||||
if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
|
||||
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
|
||||
// finishPart(req);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -295,9 +295,9 @@ public class FormingPaymentInstructionAssets implements ISessionStage {
|
|||
}
|
||||
|
||||
|
||||
sendSdfs(paymentInstructions);
|
||||
StageResult<Collection<PaymentInstruction>> stageResult = new StageResult<>(null, true);
|
||||
stageResult.setStageResult(paymentInstructions);
|
||||
Collection<SdfTable> sdfTables = sendSdfs(paymentInstructions);
|
||||
StageResult<PaymentInfo> stageResult = new StageResult<>(null, true);
|
||||
stageResult.setStageResult(new PaymentInfo(paymentInstructions, sdfTables));
|
||||
return stageResult;
|
||||
}
|
||||
|
||||
|
|
@ -307,7 +307,8 @@ public class FormingPaymentInstructionAssets implements ISessionStage {
|
|||
registryImdg.update(rgs);
|
||||
}
|
||||
|
||||
private void sendSdfs(List<PaymentInstruction> formedPaymentInstructions) {
|
||||
private Collection<SdfTable> sendSdfs(List<PaymentInstruction> formedPaymentInstructions) {
|
||||
Collection<SdfTable> sentSdfs = new ArrayList<>();
|
||||
List<SDf03> sDf03Created = new ArrayList<>();
|
||||
List<SDf12> sDf12Created = new ArrayList<>();
|
||||
for (PaymentInstruction paymentInstruction : formedPaymentInstructions) {
|
||||
|
|
@ -360,6 +361,9 @@ public class FormingPaymentInstructionAssets implements ISessionStage {
|
|||
swtExporterRequest.setType("SDF_12");
|
||||
kafkaSender.sendRequestToQueue(Consts.SWT_EXPORTER, swtExporterRequest);
|
||||
}
|
||||
if (sdf03GroupId != null) sentSdfs.add(SdfTable.SDF_03);
|
||||
if (sdf12GroupId != null) sentSdfs.add(SdfTable.SDF_12);
|
||||
return sentSdfs;
|
||||
}
|
||||
|
||||
private SDf12 newSDf12(PaymentInstruction paymentInstruction) {
|
||||
|
|
|
|||
|
|
@ -0,0 +1,31 @@
|
|||
package ru.spcex.clearing.session.stage.impl;
|
||||
|
||||
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||
import ru.spcex.platform.enumeration.SdfTable;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
|
||||
public class PaymentInfo {
|
||||
private Collection<PaymentInstruction> paymentInstructions;
|
||||
private Collection<SdfTable> sdfTypes;
|
||||
|
||||
public PaymentInfo(Collection<PaymentInstruction> paymentInstructions, SdfTable... sdfTypes) {
|
||||
this.paymentInstructions = paymentInstructions;
|
||||
this.sdfTypes = List.of(sdfTypes);
|
||||
}
|
||||
|
||||
public PaymentInfo(Collection<PaymentInstruction> paymentInstructions, Collection<SdfTable> sdfTypes) {
|
||||
this.paymentInstructions = paymentInstructions;
|
||||
this.sdfTypes = sdfTypes;
|
||||
}
|
||||
|
||||
public Collection<PaymentInstruction> getPaymentInstructions() {
|
||||
return paymentInstructions != null ? paymentInstructions : Collections.emptyList();
|
||||
}
|
||||
|
||||
public Collection<SdfTable> getSdfTypes() {
|
||||
return sdfTypes != null ? sdfTypes : Collections.emptyList();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,20 @@
|
|||
package ru.spcex.clearing.session.stage.monitor;
|
||||
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueEvent;
|
||||
|
||||
public abstract class Condition {
|
||||
|
||||
private volatile boolean metOnce = false;
|
||||
|
||||
public boolean wasMet() {
|
||||
return metOnce;
|
||||
}
|
||||
|
||||
protected void isMet() {
|
||||
this.metOnce = true;
|
||||
}
|
||||
|
||||
public abstract void event(SessionContinueEvent req);
|
||||
|
||||
public abstract String logName();
|
||||
}
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
package ru.spcex.clearing.session.stage.monitor;
|
||||
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueEvent;
|
||||
import ru.spcex.platform.enumeration.SdfTable;
|
||||
|
||||
public class SdfCondition extends Condition {
|
||||
|
||||
private final SdfTable targetSdf;
|
||||
|
||||
public SdfCondition(SdfTable targetSdf) {
|
||||
this.targetSdf = targetSdf;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void event(SessionContinueEvent req) {
|
||||
if (req.getSdfType() != null && req.getSdfType().stream().anyMatch(t -> t.equals(targetSdf))) {
|
||||
isMet();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public String logName() {
|
||||
return targetSdf.getKey();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,48 @@
|
|||
package ru.spcex.clearing.session.stage.monitor;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueEvent;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
public class SessionMonitor {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final List<Condition> conditions = new ArrayList<>();
|
||||
|
||||
private SessionMonitor() {}
|
||||
|
||||
public static SessionMonitor create() {
|
||||
return new SessionMonitor();
|
||||
}
|
||||
|
||||
public SessionMonitor addCondition(Condition condition) {
|
||||
this.conditions.add(condition);
|
||||
return this;
|
||||
}
|
||||
|
||||
public synchronized boolean isMonitorPassed(SessionContinueEvent req) {
|
||||
for (Condition condition : conditions) {
|
||||
if (condition.wasMet()) {
|
||||
continue;
|
||||
}
|
||||
condition.event(req);
|
||||
log.debug("session event {}, condition [{}] {}",
|
||||
req,
|
||||
condition.logName(),
|
||||
condition.wasMet() ? "was met" : "still waiting");
|
||||
}
|
||||
return isMonitorPassed();
|
||||
}
|
||||
|
||||
public String allConditions() {
|
||||
return conditions.stream().map(Condition::logName).collect(Collectors.joining(",", "[", "]"));
|
||||
}
|
||||
|
||||
public synchronized boolean isMonitorPassed() {
|
||||
return conditions.stream().allMatch(Condition::wasMet);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,27 @@
|
|||
package ru.spcex.clearing.session.stage.monitor;
|
||||
|
||||
import ru.spcex.platform.enumeration.SdfTable;
|
||||
import ru.spcex.platform.enumeration.Section;
|
||||
|
||||
public class SessionMonitorFactory {
|
||||
public static SessionMonitor paymentsWereCreated(Section section) {
|
||||
switch (section) {
|
||||
case MKR -> {
|
||||
return SessionMonitor.create().addCondition(new SdfCondition(SdfTable.SDF_04));
|
||||
}
|
||||
case FOND -> {
|
||||
return SessionMonitor.create()
|
||||
.addCondition(new SdfCondition(SdfTable.SDF_04))
|
||||
.addCondition(new SdfCondition(SdfTable.SDF_13));
|
||||
}
|
||||
default -> throw new IllegalStateException("unknown wait conditions for section " + section);
|
||||
}
|
||||
}
|
||||
|
||||
public static SessionMonitor waitRevise() {
|
||||
return SessionMonitor.create()
|
||||
.addCondition(new SdfCondition(SdfTable.SDF_01))
|
||||
.addCondition(new SdfCondition(SdfTable.SDF_57));
|
||||
|
||||
}
|
||||
}
|
||||
|
|
@ -4,11 +4,13 @@ import ru.spcex.platform.utils.enumeration.IEnumKey;
|
|||
|
||||
public enum SdfTable implements IEnumKey {
|
||||
SDF_01("SDF_01"),
|
||||
SDF_03("SDF_03"),
|
||||
SDF_04("SDF_04"),
|
||||
SDF_06("SDF_06"),
|
||||
SDF_08("SDF_08"),
|
||||
SDF_09("SDF_09"),
|
||||
SDF_10("SDF_10"),
|
||||
SDF_12("SDF_12"),
|
||||
SDF_13("SDF_13"),
|
||||
SDF_16("SDF_16"),
|
||||
SDF_21("SDF_21"),
|
||||
|
|
|
|||
|
|
@ -1,16 +0,0 @@
|
|||
package ru.spcex.clearing.platform.messaging.domain.cud.clearing;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
|
||||
public class ContinueSessionBnRequest {
|
||||
@JsonProperty
|
||||
private Long generationId;
|
||||
|
||||
public Long getGenerationId() {
|
||||
return generationId;
|
||||
}
|
||||
|
||||
public void setGenerationId(Long generationId) {
|
||||
this.generationId = generationId;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,34 @@
|
|||
package ru.spcex.clearing.platform.messaging.domain.cud.clearing;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
import ru.spcex.platform.enumeration.SdfTable;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
public class SessionContinueEvent {
|
||||
|
||||
public SessionContinueEvent(SdfTable... sdfType) {
|
||||
this.sdfTypes = List.of(sdfType);
|
||||
}
|
||||
|
||||
public SessionContinueEvent() {
|
||||
}
|
||||
|
||||
@JsonProperty
|
||||
private List<SdfTable> sdfTypes;
|
||||
|
||||
public List<SdfTable> getSdfType() {
|
||||
return sdfTypes;
|
||||
}
|
||||
|
||||
public void setSdfType(List<SdfTable> sdfType) {
|
||||
this.sdfTypes = sdfType;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "SessionContinueEvent{" +
|
||||
"sdfType=" + sdfTypes +
|
||||
'}';
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue