diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java index 56c575ae9..c189a6cf4 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/SdfGroup.java @@ -1,14 +1,13 @@ package ru.spcex.clearing.statement; -import ru.spcex.platform.enumeration.SdfTable; - import java.util.Arrays; import java.util.Collection; import java.util.List; import java.util.Optional; +import ru.spcex.platform.enumeration.SdfTable; public enum SdfGroup { - Sdf01And57(SdfTable.SDF_01, SdfTable.SDF_57), +// Sdf01And57(SdfTable.SDF_01, SdfTable.SDF_57), Sdf08And21(SdfTable.SDF_08, SdfTable.SDF_21), //------ session groups ------ (4), (1 57), (13), (8 21) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java index 5fc3a8031..29b9aeb53 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/statement/StatementServiceV2.java @@ -1,10 +1,27 @@ package ru.spcex.clearing.statement; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.Iterator; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.function.Predicate; +import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; import ru.clearing.classes.statics.data.misc.Session; -import ru.clearing.classes.statics.data.sdf.*; +import ru.clearing.classes.statics.data.sdf.SDf01; +import ru.clearing.classes.statics.data.sdf.SDf04; +import ru.clearing.classes.statics.data.sdf.SDf08; +import ru.clearing.classes.statics.data.sdf.SDf13; +import ru.clearing.classes.statics.data.sdf.SDf21; +import ru.clearing.classes.statics.data.sdf.SDf57; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; @@ -14,7 +31,17 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueE import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.service.StatementService; -import ru.spcex.clearing.service.executors.*; +import ru.spcex.clearing.service.executors.AbstractExecutor; +import ru.spcex.clearing.service.executors.Reviser; +import ru.spcex.clearing.service.executors.Sdf01Executor; +import ru.spcex.clearing.service.executors.Sdf04Executor; +import ru.spcex.clearing.service.executors.Sdf08Executor; +import ru.spcex.clearing.service.executors.Sdf10Executor; +import ru.spcex.clearing.service.executors.Sdf13Executor; +import ru.spcex.clearing.service.executors.Sdf20Executor; +import ru.spcex.clearing.service.executors.Sdf21Executor; +import ru.spcex.clearing.service.executors.Sdf55Executor; +import ru.spcex.clearing.service.executors.Sdf57Executor; import ru.spcex.clearing.service.model.Result; import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.enumeration.ObjectType; @@ -25,10 +52,6 @@ import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.text.TextUtil; -import java.util.*; -import java.util.function.Predicate; -import java.util.stream.Collectors; - @Component public class StatementServiceV2 { private Logger log = LoggerFactory.getLogger(getClass()); @@ -93,11 +116,13 @@ public class StatementServiceV2 { log.debug("table {} is not paired with any other SDF", table); //здесь switch case одиночные методы switch (table) { + case SDF_01 -> processSdf01(systemRequest.getRequestPayload()); case SDF_04 -> processSdf04(systemRequest.getRequestPayload()); case SDF_10 -> sdf10Executor.execute(systemRequest); case SDF_13 -> processSdf13(systemRequest.getRequestPayload()); case SDF_20 -> sdf20Executor.execute(systemRequest); case SDF_55 -> sdf55Executor.execute(systemRequest); + case SDF_57 -> processSdf57(systemRequest.getRequestPayload()); default -> log.error("unknown table {}", table); } return; @@ -125,9 +150,10 @@ public class StatementServiceV2 { .collect(TextUtil.join)); if (sdfGroup.get() == SdfGroup.Sdf08And21) { processSdf08And21(find(SdfTable.SDF_08, fullGroup), find(SdfTable.SDF_21, fullGroup)); - } else if (sdfGroup.get() == SdfGroup.Sdf01And57) { - processSdf01And57(find(SdfTable.SDF_01, fullGroup), find(SdfTable.SDF_57, fullGroup)); } +// else if (sdfGroup.get() == SdfGroup.Sdf01And57) { +// processSdf01Parent(find(SdfTable.SDF_01, fullGroup), find(SdfTable.SDF_57, fullGroup)); +// } //else if (sdfGroup.get() == SdfGroup.session_Triple) { // processSdf04(find(SdfTable.SDF_04, fullGroup)); // processSdf01And57(find(SdfTable.SDF_01, fullGroup), find(SdfTable.SDF_57, fullGroup)); @@ -172,32 +198,6 @@ public class StatementServiceV2 { log.info("pair sdf08/sdf21 processed successfully"); } - /** - * fromAccService передается когда пришел ответ от account-service - * в этом случае: по key находим пару в которой сохранен sdf57 запрос и частично выполненный sdf01 - * вместо старого sdf01 запроса выполняем новый пришедший от account-service - */ - private void processSdf01And57(StatementRequest sdf01, StatementRequest sdf57) { - Result sdf01Res = processSdf01(sdf01); - removeFirstWithSameTableAndGroupId(sdf01); - if (sdf01Res.getAccountRequests().size() > 0) { - log.info("sdf01 execution wasn't complete, waiting for an answer from account-service"); - return; - } - //затем sdf57 - boolean sessionStarted = processSdf57(sdf57); - removeFirstWithSameTableAndGroupId(sdf57); - reviser.doRevise(sdf01.getGroupId()); - //теперь можем продолжить сессию с шага 1 - if (!sessionStarted - || sessionImdg.getFirstObjectBySQL("workflowStatus = '%s'".formatted(WorkflowStatus.Active.getKey())) != null) { - SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_01, SdfTable.SDF_57); - kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn); - } - log.info("pair sdf01/sdf57 processed successfully"); - - } - private void processSdf04(StatementRequest statementRequest) { Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class); Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( @@ -221,41 +221,52 @@ public class StatementServiceV2 { } } - private Result processSdf01(StatementRequest statementRequest) { + private void processSdf01(StatementRequest sdf01) { Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class); Collection sdfGroup; - if (statementRequest.getAccountCreationResults().size() == 0) { + if (sdf01.getAccountCreationResults().size() == 0) { sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( - "generationId", statementRequest.getGroupId())); + "generationId", sdf01.getGroupId())); } else { log.debug("accounts were created: {} for generationId {} child generationId {}", - statementRequest.getAccountCreationResults().size(), statementRequest.getGroupId(), statementRequest.getChildGenerationId()); + sdf01.getAccountCreationResults().size(), sdf01.getGroupId(), sdf01.getChildGenerationId()); //todo могут ли тут реквесты бегать по кругу, если да, то у нас проблемы - sdfGroup = statementRequest.getAccountCreationResults() + sdfGroup = sdf01.getAccountCreationResults() .stream() // .filter(part -> part.getErrorCode() == null) //fixme эти случае должны попадать в ошибочный sdf02 .map(part -> sdfImdg.getSingleObjectByID(part.getSdfId())) .collect(Collectors.toList()); } AbstractExecutor service = sdf01Executor; - Result res = service.execute(sdfGroup, statementRequest); - if (res.getAccountRequests().size() != 0) { - kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, StatementService.createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests(), res.getChildGenerationId())); - return res; + Result sdf01Res = service.execute(sdfGroup, sdf01); + if (sdf01Res.getAccountRequests().size() != 0) { + kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, StatementService.createAccountsRequest(sdf01.getGroupId(), sdf01Res.getAccountRequests(), sdf01Res.getChildGenerationId())); } else if (service.isNeedToSendCommand()) { - service.sendCommand(kafkaSender, res); - return res; + service.sendCommand(kafkaSender, sdf01Res); } - return res; + removeFirstWithSameTableAndGroupId(sdf01); + if (!sdf01Res.getAccountRequests().isEmpty()) { + log.info("sdf01 execution wasn't complete, waiting for an answer from account-service"); + return; + } + reviser.doRevise(sdf01.getGroupId()); + log.info("sdf01 processed successfully"); } - private boolean processSdf57(StatementRequest statementRequest) { + private void processSdf57(StatementRequest sdf57) { Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class); Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( - "generationId", statementRequest.getGroupId())); + "generationId", sdf57.getGroupId())); AbstractExecutor service = sdf57Executor; - Result res = service.execute(sdfGroup, statementRequest); - return res.isSessionWasStarted(); + Result res = service.execute(sdfGroup, sdf57); + boolean sessionStarted = res.isSessionWasStarted(); + removeFirstWithSameTableAndGroupId(sdf57); + if (!sessionStarted + || sessionImdg.getFirstObjectBySQL("workflowStatus = '%s'".formatted(WorkflowStatus.Active.getKey())) != null) { + SessionContinueEvent continueSessionBn = new SessionContinueEvent(SdfTable.SDF_01, SdfTable.SDF_57); + kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn); + } + log.info("sdf57 processed successfully"); // finishSendCommand(res, service, statementRequest); }