session sdf04/sdf13
This commit is contained in:
parent
e57ee9a72e
commit
4316830043
7 changed files with 218 additions and 92 deletions
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<ExecutionDeposit> executionDepositImdg;
|
||||
private final Supplier<List<String>> marketCodes;
|
||||
private final Imdg<Registry> 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<List<ExecutionCommon>> 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) {
|
||||
|
|
|
|||
|
|
@ -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<ExecutionFond> executionFondImdg;
|
||||
private final Supplier<List<String>> 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<List<ExecutionCommon>> dealsPreparationResult;
|
||||
{
|
||||
|
|
@ -164,26 +182,19 @@ public class PrimaryAuctionB0Session extends AbstractSession implements Initiali
|
|||
//stage 7
|
||||
paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstruction);
|
||||
}
|
||||
//stage 7
|
||||
// StageResult<Collection<PaymentInstruction>> 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) {
|
||||
|
|
|
|||
|
|
@ -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<ExecutionFond> executionFondImdg;
|
||||
private final Supplier<List<String>> 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<List<ExecutionCommon>> dealsPreparationResult;
|
||||
{
|
||||
|
|
@ -166,26 +184,19 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali
|
|||
//stage 7
|
||||
paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets);
|
||||
}
|
||||
//stage 7
|
||||
// StageResult<Collection<PaymentInstruction>> 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) {
|
||||
|
|
|
|||
|
|
@ -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<ExecutionFond> executionFondImdg;
|
||||
private final Supplier<List<String>> 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<List<ExecutionCommon>> dealsPreparationResult;
|
||||
{
|
||||
|
|
@ -165,26 +183,19 @@ public class PrimaryAuctionT0Session extends AbstractSession implements Initiali
|
|||
//stage 7
|
||||
paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets);
|
||||
}
|
||||
//stage 7
|
||||
// StageResult<Collection<PaymentInstruction>> 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) {
|
||||
|
|
|
|||
|
|
@ -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<Registry> 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) {
|
||||
|
|
|
|||
|
|
@ -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<ExecutionFond> executionFondImdg;
|
||||
private final Supplier<List<String>> 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<List<ExecutionCommon>> 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) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue