From 431683004303905d0053489a8009fee05f9f3103 Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 8 Aug 2023 20:45:56 +0300 Subject: [PATCH] session sdf04/sdf13 --- .../clearing/service/StatementService.java | 2 +- .../session/stage/IntermediateMkrSession.java | 46 +++++++++++---- .../stage/PrimaryAuctionB0Session.java | 56 ++++++++++++------- .../stage/PrimaryAuctionBnSession.java | 56 ++++++++++++------- .../stage/PrimaryAuctionT0Session.java | 56 ++++++++++++------- .../session/stage/ReturnDepositSession.java | 48 ++++++++++++---- .../stage/SecondaryAuctionT0Session.java | 46 +++++++++++---- 7 files changed, 218 insertions(+), 92 deletions(-) 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 7ff36b09e..46a090606 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 @@ -71,7 +71,7 @@ public class StatementService extends QueueConsumer implements InitializingBean .setConsumer(systemRequest -> { StatementRequest payload = systemRequest.getRequestPayload(); if (payload.getTable() == null - || !Arrays.asList(SdfTable.SDF_01, SdfTable.SDF_57, SdfTable.SDF_08, SdfTable.SDF_13, SdfTable.SDF_21).contains(payload.getTable())) { + || !Arrays.asList(SdfTable.SDF_01, SdfTable.SDF_04, SdfTable.SDF_57, SdfTable.SDF_08, SdfTable.SDF_13, SdfTable.SDF_21).contains(payload.getTable())) { log.debug("StatementService: skipping table {}", payload.getTable()); return; } 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 0674d986e..309786f99 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 @@ -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,6 +53,9 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ private final Imdg executionDepositImdg; private final Supplier> marketCodes; private final Imdg registryImdg; + private SessionMonitor firstReviseMonitor; + private SessionMonitor afterPaymentsSdf4Monitor; + private SessionMonitor afterPaymentsReviseMonitor; public IntermediateMkrSession( @@ -120,6 +125,7 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ endSession(); } else { log.info("stage BalanceRevise success, waiting for a response from kafka"); + this.firstReviseMonitor = SessionMonitorFactory.waitRevise(); } } @@ -128,20 +134,32 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ 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; { @@ -193,15 +211,18 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDeals); } if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) { - runStage(TaskType.FormingPaymentInstruction, balanceRevise); -// finishPart(req); + 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) { + public void finishPart() { try { if (!checkStage(TaskType.FormingPaymentInstruction)) { log.error("cannot continue session, current stage is {}", currStage.get()); @@ -227,6 +248,11 @@ public class IntermediateMkrSession extends AbstractSession implements Initializ } } + private void sendSdf56() { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + } + private boolean startSession() { synchronized (this.currStage) { if (this.currStage.get() != null) { 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 4b3aa0eb9..98ad5f1b8 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 @@ -11,6 +11,8 @@ 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.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.Section; @@ -46,6 +48,9 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali private final Imdg executionFondImdg; private final Supplier> marketCodes; + private SessionMonitor firstReviseMonitor; + private SessionMonitor afterPaymentsSdf4And13Monitor; + private SessionMonitor afterPaymentsReviseMonitor; public PrimaryAuctionB0Session( @@ -104,6 +109,7 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali endSession(); } else { log.info("stage BalanceRevise success, waiting for a response from kafka"); + this.firstReviseMonitor = SessionMonitorFactory.waitRevise(); } } @@ -112,20 +118,32 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali 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, afterPaymentsSdf4And13Monitor, afterPaymentsReviseMonitor)); + if (firstReviseMonitor != null && isMonitorPassed(firstReviseMonitor, req.getRequestPayload())) { + firstReviseMonitor = null; + firstPart(); + } + if (afterPaymentsSdf4And13Monitor != null && isMonitorPassed(afterPaymentsSdf4And13Monitor, req.getRequestPayload())) { + afterPaymentsSdf4And13Monitor = 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; { @@ -164,26 +182,19 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali //stage 7 paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstruction); } - //stage 7 -// StageResult> paymentResult = null; -// { -// FormingPaymentInstructionDealsMkrPayload payload = new FormingPaymentInstructionDealsMkrPayload(); -// payload.setSessionId(currSession.getId()); -// payload.setSection(section()); -// payload.setPaymentInstructionReturns(returnsPayment.getStageResult()); -// -// paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDealsFinalMkr); -// } if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) { - runStage(TaskType.FormingPaymentInstruction, balanceRevise); -// finishPart(req); + log.info("no payment instructions were created, sending SDF56"); + sendSdf56(); + } else { + log.info("created {} PaymentInstructions, waiting for SDF04/SDF13", paymentResult.getStageResult().getPaymentInstructions().size()); + this.afterPaymentsSdf4And13Monitor = SessionMonitorFactory.paymentsWereCreated(section()); } } catch (StageException e) { //already logged } } - public void finishPart(BaseRequest req) { + public void finishPart() { try { if (!checkStage(TaskType.FormingPaymentInstruction)) { log.error("cannot continue session, current stage is {}", currStage.get()); @@ -209,6 +220,11 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali } } + private void sendSdf56() { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + } + private boolean startSession() { synchronized (this.currStage) { if (this.currStage.get() != null) { 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 d227d0776..efd9df20a 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 @@ -11,6 +11,8 @@ 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.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.Section; @@ -46,6 +48,9 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali private final Imdg executionFondImdg; private final Supplier> marketCodes; + private SessionMonitor firstReviseMonitor; + private SessionMonitor afterPaymentsSdf4And13Monitor; + private SessionMonitor afterPaymentsReviseMonitor; public PrimaryAuctionBnSession( @@ -105,6 +110,7 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali endSession(); } else { log.info("stage BalanceRevise success, waiting for a response from kafka"); + this.firstReviseMonitor = SessionMonitorFactory.waitRevise(); } } @@ -113,20 +119,32 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali 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, afterPaymentsSdf4And13Monitor, afterPaymentsReviseMonitor)); + if (firstReviseMonitor != null && isMonitorPassed(firstReviseMonitor, req.getRequestPayload())) { + firstReviseMonitor = null; + firstPart(); + } + if (afterPaymentsSdf4And13Monitor != null && isMonitorPassed(afterPaymentsSdf4And13Monitor, req.getRequestPayload())) { + afterPaymentsSdf4And13Monitor = 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; { @@ -166,26 +184,19 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali //stage 7 paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets); } - //stage 7 -// StageResult> paymentResult = null; -// { -// FormingPaymentInstructionDealsMkrPayload payload = new FormingPaymentInstructionDealsMkrPayload(); -// payload.setSessionId(currSession.getId()); -// payload.setSection(section()); -// payload.setPaymentInstructionReturns(returnsPayment.getStageResult()); -// -// paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDealsFinalMkr); -// } if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) { - runStage(TaskType.FormingPaymentInstruction, balanceRevise); -// finishPart(req); + log.info("no payment instructions were created, sending SDF56"); + sendSdf56(); + } else { + log.info("created {} PaymentInstructions, waiting for SDF04/SDF13", paymentResult.getStageResult().getPaymentInstructions().size()); + this.afterPaymentsSdf4And13Monitor = SessionMonitorFactory.paymentsWereCreated(section()); } } catch (StageException e) { //already logged } } - public void finishPart(BaseRequest req) { + public void finishPart() { try { if (!checkStage(TaskType.FormingPaymentInstruction)) { log.error("cannot continue session, current stage is {}", currStage.get()); @@ -211,6 +222,11 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali } } + private void sendSdf56() { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + } + private boolean startSession() { synchronized (this.currStage) { if (this.currStage.get() != null) { 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 ab7188ab8..0c3f093a2 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 @@ -11,6 +11,8 @@ 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.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.Section; @@ -46,6 +48,9 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali private final Imdg executionFondImdg; private final Supplier> marketCodes; + private SessionMonitor firstReviseMonitor; + private SessionMonitor afterPaymentsSdf4And13Monitor; + private SessionMonitor afterPaymentsReviseMonitor; public PrimaryAuctionT0Session( ImdgProvider imdgProvider, @@ -105,6 +110,7 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali endSession(); } else { log.info("stage BalanceRevise success, waiting for a response from kafka"); + this.firstReviseMonitor = SessionMonitorFactory.waitRevise(); } } @@ -113,20 +119,32 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali 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, afterPaymentsSdf4And13Monitor, afterPaymentsReviseMonitor)); + if (firstReviseMonitor != null && isMonitorPassed(firstReviseMonitor, req.getRequestPayload())) { + firstReviseMonitor = null; + firstPart(); + } + if (afterPaymentsSdf4And13Monitor != null && isMonitorPassed(afterPaymentsSdf4And13Monitor, req.getRequestPayload())) { + afterPaymentsSdf4And13Monitor = 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; { @@ -165,26 +183,19 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali //stage 7 paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets); } - //stage 7 -// StageResult> paymentResult = null; -// { -// FormingPaymentInstructionDealsMkrPayload payload = new FormingPaymentInstructionDealsMkrPayload(); -// payload.setSessionId(currSession.getId()); -// payload.setSection(section()); -// payload.setPaymentInstructionReturns(returnsPayment.getStageResult()); -// -// paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionDealsFinalMkr); -// } if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) { - runStage(TaskType.FormingPaymentInstruction, balanceRevise); -// finishPart(req); + log.info("no payment instructions were created, sending SDF56"); + sendSdf56(); + } else { + log.info("created {} PaymentInstructions, waiting for SDF04/SDF13", paymentResult.getStageResult().getPaymentInstructions().size()); + this.afterPaymentsSdf4And13Monitor = SessionMonitorFactory.paymentsWereCreated(section()); } } catch (StageException e) { //already logged } } - public void finishPart(BaseRequest req) { + public void finishPart() { try { if (!checkStage(TaskType.FormingPaymentInstruction)) { log.error("cannot continue session, current stage is {}", currStage.get()); @@ -210,6 +221,11 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali } } + private void sendSdf56() { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + } + private boolean startSession() { synchronized (this.currStage) { if (this.currStage.get() != null) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java index c6d6c206c..5342799de 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/ReturnDepositSession.java @@ -10,6 +10,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.enumeration.RegistryStatus; import ru.spcex.platform.enumeration.Section; @@ -40,6 +42,9 @@ public class ReturnDepositSession extends AbstractSession implements Initializin private final FinishingSession finishingSession; private final EndStageNotification endStageNotification; private final Imdg registryImdg; + private SessionMonitor firstReviseMonitor; + private SessionMonitor afterPaymentsSdf4Monitor; + private SessionMonitor afterPaymentsReviseMonitor; public ReturnDepositSession( @@ -99,6 +104,7 @@ public class ReturnDepositSession extends AbstractSession implements Initializin endSession(); } else { log.info("stage BalanceRevise success, waiting for a response from kafka"); + this.firstReviseMonitor = SessionMonitorFactory.waitRevise(); } } @@ -107,20 +113,32 @@ public class ReturnDepositSession extends AbstractSession implements Initializin 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 4 { InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload(); @@ -147,16 +165,19 @@ public class ReturnDepositSession extends AbstractSession implements Initializin //stage 7 paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstruction); } - if (paymentResult != null && paymentResult.getStageResult().isEmpty()) { - runStage(TaskType.FormingPaymentInstruction, balanceRevise); -// finishPart(req); + if (paymentResult == null || paymentResult.getStageResult().isEmpty()) { + log.info("no payment instructions were created, sending SDF56"); + sendSdf56(); + } else { + log.info("created {} PaymentInstructions, waiting for SDF04", paymentResult.getStageResult().size()); + this.afterPaymentsSdf4Monitor = SessionMonitorFactory.paymentsWereCreated(section()); } } catch (StageException e) { //already logged } } - public void finishPart(BaseRequest req) { + public void finishPart() { try { if (!checkStage(TaskType.FormingPaymentInstruction)) { log.error("cannot continue session, current stage is {}", currStage.get()); @@ -181,6 +202,11 @@ public class ReturnDepositSession extends AbstractSession implements Initializin } } + private void sendSdf56() { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + } + private boolean startSession() { synchronized (this.currStage) { if (this.currStage.get() != null) { 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 963984751..b6a198499 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 @@ -11,6 +11,8 @@ 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.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.Section; @@ -45,6 +47,9 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia private final Imdg executionFondImdg; private final Supplier> marketCodes; + private SessionMonitor firstReviseMonitor; + private SessionMonitor afterPaymentsSdf4And13Monitor; + private SessionMonitor afterPaymentsReviseMonitor; public SecondaryAuctionT0Session( @@ -101,6 +106,7 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia endSession(); } else { log.info("stage BalanceRevise success, waiting for a response from kafka"); + this.firstReviseMonitor = SessionMonitorFactory.waitRevise(); } } @@ -109,20 +115,32 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia 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, afterPaymentsSdf4And13Monitor, afterPaymentsReviseMonitor)); + if (firstReviseMonitor != null && isMonitorPassed(firstReviseMonitor, req.getRequestPayload())) { + firstReviseMonitor = null; + firstPart(); + } + if (afterPaymentsSdf4And13Monitor != null && isMonitorPassed(afterPaymentsSdf4And13Monitor, req.getRequestPayload())) { + afterPaymentsSdf4And13Monitor = 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; { @@ -163,15 +181,18 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia } if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) { - runStage(TaskType.FormingPaymentInstruction, balanceRevise); -// finishPart(req); + log.info("no payment instructions were created, sending SDF56"); + sendSdf56(); + } else { + log.info("created {} PaymentInstructions, waiting for SDF04/SDF13", paymentResult.getStageResult().getPaymentInstructions().size()); + this.afterPaymentsSdf4And13Monitor = SessionMonitorFactory.paymentsWereCreated(section()); } } catch (StageException e) { //already logged } } - public void finishPart(BaseRequest req) { + public void finishPart() { try { if (!checkStage(TaskType.FormingPaymentInstruction)) { log.error("cannot continue session, current stage is {}", currStage.get()); @@ -197,6 +218,11 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia } } + private void sendSdf56() { + runStage(TaskType.FormingPaymentInstruction, balanceRevise); + this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise(); + } + private boolean startSession() { synchronized (this.currStage) { if (this.currStage.get() != null) {