split sdf57/sdf01
This commit is contained in:
parent
9971299c70
commit
b004b9191a
2 changed files with 63 additions and 53 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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<SDf04> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class);
|
||||
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
|
||||
|
|
@ -221,41 +221,52 @@ public class StatementServiceV2 {
|
|||
}
|
||||
}
|
||||
|
||||
private Result processSdf01(StatementRequest statementRequest) {
|
||||
private void processSdf01(StatementRequest sdf01) {
|
||||
Imdg<SDf01> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
|
||||
Collection<? extends SpcexObjectBase> 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<SDf57> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class);
|
||||
Collection<? extends SpcexObjectBase> 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);
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue