diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index 31ea1eaf3..1e794c944 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -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); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index 724227a21..fc814b327 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -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); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java index b96295d77..2e382d899 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf57Executor.java @@ -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 { @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 diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java index 390714828..8171341d3 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/AbstractSession.java @@ -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 ""; + return Arrays.stream(monitors) + .filter(Objects::nonNull) + .map(SessionMonitor::allConditions) + .collect(Collectors.joining("/")); + } + protected StageResult runStage(TaskType type, ISessionStage stage) { return runStage(type, null, stage); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java index cb2f2e7ee..2a73e40dc 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/FinalMkrSession.java @@ -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 executionDepositImdg; private final Supplier> marketCodes; private final Imdg 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> 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 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> returnsPayment = null; + StageResult> 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> paymentResult = null; + StageResult 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()); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java index b3a6b4e49..0674d986e 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/IntermediateMkrSession.java @@ -184,7 +184,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ returnsPayment = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionReturn); } - StageResult> paymentResult = null; + StageResult 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); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java index c12554259..4b3aa0eb9 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionB0Session.java @@ -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 } - StageResult> paymentResult = null; + StageResult 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); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java index c6220873a..d227d0776 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java @@ -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 } //stage 7 - StageResult> paymentResult = null; + StageResult 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); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java index 1847dae68..ab7188ab8 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionT0Session.java @@ -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 } - StageResult> paymentResult = null; + StageResult 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); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java index 2e3c2340f..963984751 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SecondaryAuctionT0Session.java @@ -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 } //stage 7 - StageResult> paymentResult = null; + StageResult 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); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java index 26c9ce795..54a9aaa4a 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstructionAssets.java @@ -295,9 +295,9 @@ public class FormingPaymentInstructionAssets implements ISessionStage { } - sendSdfs(paymentInstructions); - StageResult> stageResult = new StageResult<>(null, true); - stageResult.setStageResult(paymentInstructions); + Collection sdfTables = sendSdfs(paymentInstructions); + StageResult 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 formedPaymentInstructions) { + private Collection sendSdfs(List formedPaymentInstructions) { + Collection sentSdfs = new ArrayList<>(); List sDf03Created = new ArrayList<>(); List 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) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/PaymentInfo.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/PaymentInfo.java new file mode 100644 index 000000000..c791ce82b --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/PaymentInfo.java @@ -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 paymentInstructions; + private Collection sdfTypes; + + public PaymentInfo(Collection paymentInstructions, SdfTable... sdfTypes) { + this.paymentInstructions = paymentInstructions; + this.sdfTypes = List.of(sdfTypes); + } + + public PaymentInfo(Collection paymentInstructions, Collection sdfTypes) { + this.paymentInstructions = paymentInstructions; + this.sdfTypes = sdfTypes; + } + + public Collection getPaymentInstructions() { + return paymentInstructions != null ? paymentInstructions : Collections.emptyList(); + } + + public Collection getSdfTypes() { + return sdfTypes != null ? sdfTypes : Collections.emptyList(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/Condition.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/Condition.java new file mode 100644 index 000000000..a6325e0eb --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/Condition.java @@ -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(); +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SdfCondition.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SdfCondition.java new file mode 100644 index 000000000..c7deb6a51 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SdfCondition.java @@ -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(); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SessionMonitor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SessionMonitor.java new file mode 100644 index 000000000..e914561e0 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SessionMonitor.java @@ -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 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); + } + +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SessionMonitorFactory.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SessionMonitorFactory.java new file mode 100644 index 000000000..b391c93ee --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/monitor/SessionMonitorFactory.java @@ -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)); + + } +} diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java index 7edf8fbab..e83326702 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java @@ -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"), diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java deleted file mode 100644 index e02afec3c..000000000 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/ContinueSessionBnRequest.java +++ /dev/null @@ -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; - } -} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/SessionContinueEvent.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/SessionContinueEvent.java new file mode 100644 index 000000000..b6c8fbd7f --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/SessionContinueEvent.java @@ -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 sdfTypes; + + public List getSdfType() { + return sdfTypes; + } + + public void setSdfType(List sdfType) { + this.sdfTypes = sdfType; + } + + @Override + public String toString() { + return "SessionContinueEvent{" + + "sdfType=" + sdfTypes + + '}'; + } +}