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 fc814b327..b672520ca 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 @@ -38,6 +38,7 @@ public class StatementService extends QueueConsumer implements InitializingBean private final Map> sdfImdgs; private final Map> executorsMap; private final Reviser reviser; + private final StatementServiceV2 stmtSrvV2; /** * Пары sdf запросов пришедшие с модуля dbf-import @@ -48,10 +49,11 @@ public class StatementService extends QueueConsumer implements InitializingBean public StatementService(Consumer kafkaQueue, ImdgProvider imdgProvider, KafkaSender kafkaSender, - @Qualifier("sdfExecutors") Map> executorsMap, Reviser reviser) { + @Qualifier("sdfExecutors") Map> executorsMap, Reviser reviser, StatementServiceV2 stmtSrvV2) { super(kafkaQueue); this.imdgProvider = imdgProvider; this.reviser = reviser; + this.stmtSrvV2 = stmtSrvV2; this.sdfImdgs = new EnumMap<>(SdfTable.class); this.kafkaSender = kafkaSender; this.executorsMap = executorsMap; @@ -68,7 +70,9 @@ public class StatementService extends QueueConsumer implements InitializingBean callback(StatementRequest.class) .setConsumer(systemRequest -> { StatementRequest payload = systemRequest.getRequestPayload(); - if (payload.isContinueSdf()) { + if (Arrays.asList(SdfTable.SDF_08, SdfTable.SDF_21).contains(payload.getTable())) { + stmtSrvV2.processReq(systemRequest); + } else if (payload.isContinueSdf()) { processAccountAnswer(systemRequest); } else { processPaired(systemRequest); @@ -382,7 +386,7 @@ public class StatementService extends QueueConsumer implements InitializingBean return r; } - private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List accountRequests, Long childGenerationId) { + public static AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List accountRequests, Long childGenerationId) { AccountSdf01Request r = new AccountSdf01Request(); r.setGroupingSdf01Id(sdf01GroupingId); r.setGroupingSdf02Id(childGenerationId); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementServiceV2.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementServiceV2.java new file mode 100644 index 000000000..3876a9781 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementServiceV2.java @@ -0,0 +1,169 @@ +package ru.spcex.clearing.service; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.sdf.SDf08; +import ru.clearing.classes.statics.data.sdf.SDf21; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +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.balance.StatementRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.service.executors.Sdf08Executor; +import ru.spcex.clearing.service.executors.Sdf21Executor; +import ru.spcex.clearing.service.model.Result; +import ru.spcex.platform.enumeration.SdfTable; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.*; +import java.util.stream.Collectors; + +@Component +public class StatementServiceV2 { + private Logger log = LoggerFactory.getLogger(getClass()); + private final ImdgProvider imdgProvider; + private final KafkaSender kafkaSender; + private final List statementRequests = new LinkedList<>(); + private final Sdf08Executor sdf08Executor; + private final Sdf21Executor sdf21Executor; + + public StatementServiceV2(ImdgProvider imdgProvider, KafkaSender kafkaSender, Sdf08Executor sdf08Executor, Sdf21Executor sdf21Executor) { + this.imdgProvider = imdgProvider; + this.kafkaSender = kafkaSender; + this.sdf08Executor = sdf08Executor; + this.sdf21Executor = sdf21Executor; + } + + public void processReq(BaseRequest systemRequest) { + SdfTable table = systemRequest.getRequestPayload().getTable(); + Long groupId = systemRequest.getRequestPayload().getGroupId(); + log.debug("received stmtReq.SdfTable={}, groupId={}", table, groupId); + Optional sdfGroup = SdfGroup.groupByTable(table); + if (sdfGroup.isEmpty()) { + log.debug("table {} is not paired with any other SDF", table); + //здесь switch case одиночные методы + return; + } + Collection fullGroup = getFullGroup(systemRequest.getRequestPayload(), sdfGroup.get()); + if (fullGroup.isEmpty()) { + log.debug("didn't find full sdf set for stmtReq.SdfTable={}, groupId={}. Adding to cache", table, groupId); + statementRequests.add(systemRequest.getRequestPayload()); + return; + } + log.debug("find full set of SDF requests: {}", fullGroup.stream() + .map(stReq -> stReq.getTable().getKey() + " generationId=" + stReq.getGroupId()) + .collect(Collectors.joining(",", "[", "]"))); + if (sdfGroup.get() == SdfGroup.Sdf08And21) { + processSdf08And21(find(SdfTable.SDF_08, fullGroup), find(SdfTable.SDF_21, fullGroup)); + } else { + throw new IllegalStateException("not implemented"); + } + } + + private static StatementRequest find(SdfTable table, Collection reqs) { + Optional first = reqs.stream().filter(r -> r.getTable().equals(table)).findFirst(); + if (first.isEmpty()) { + throw new IllegalStateException("couldn't find + " + table + " - this should have happened "); + } + return first.get(); + } + + /** + * fromAccService передается когда пришел ответ от account-service + * в этом случае: по key находим пару в которой сохранен sdf57 запрос и частично выполненный sdf01 + * вместо старого sdf01 запроса выполняем новый пришедший от account-service + */ + private void processSdf08And21(StatementRequest sdf08, StatementRequest sdf21) { + Result sdf08Res = processSdf08(sdf08); + if (sdf08Res.getAccountRequests().size() > 0) { + statementRequests.remove(sdf08); + log.info("sdf08 execution wasn't complete, waiting for an answer from account-service"); + return; + } + //затем sdf21 + processSdf21(sdf21); + //fixme ревизия для бумаг reviser.doRevise(pair.getFirst().getGroupId()); + log.info("pair sdf08/sdf21 processed successfully"); + statementRequests.removeAll(Arrays.asList(sdf08, sdf21)); + } + + + private Result processSdf08(StatementRequest statementRequest) { + Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class); + Collection sdfGroup; + if (statementRequest.getAccountCreationResults().size() == 0) { + sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( + "generationId", statementRequest.getGroupId())); + } else { + sdfGroup = statementRequest.getAccountCreationResults() + .stream() + .filter(part -> part.getErrorCode() == null) + .map(part -> sdfImdg.getSingleObjectByID(part.getSdfId())) + .collect(Collectors.toList()); + } + Result res = sdf08Executor.execute(sdfGroup, statementRequest); + if (res.getAccountRequests().size() != 0) { + AccountSdf01Request createAccsReq = StatementService.createAccountsRequest(statementRequest.getGroupId(), + res.getAccountRequests(), + res.getGenerationId()); + kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF08, createAccsReq); + } else if (sdf08Executor.isNeedToSendCommand()) { + sdf08Executor.sendCommand(kafkaSender, res); + } + return res; + } + + private void processSdf21(StatementRequest statementRequest) { + Imdg sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf21, SDf21.class); + Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( + "generationId", statementRequest.getGroupId())); + Result res = sdf21Executor.execute(sdfGroup, statementRequest); + } + + + private Collection getFullGroup(StatementRequest stmtReq, SdfGroup group) { + List reqs = new ArrayList<>(); + log.debug("trying to find request pair for {}. cashed reqs {}", stmtReq.getTable(), + statementRequests.stream() + .map(StatementRequest::getTable) + .map(Object::toString) + .collect(Collectors.joining(",", "[", "]"))); + for (StatementRequest stmtReqOld : statementRequests) { //проходим с головы (с самых старых) + boolean matchGroup = group.getGroup().contains(stmtReqOld.getTable()); + boolean notIncludedThisTableAlready = reqs.stream().map(StatementRequest::getTable).noneMatch(t -> t.equals(stmtReqOld.getTable())); + boolean isNotTheSameAsNewReq = Objects.equals(stmtReqOld.getTable(), stmtReq.getTable()); + if (matchGroup && notIncludedThisTableAlready && isNotTheSameAsNewReq) { + reqs.add(stmtReqOld); + } + } + reqs.add(stmtReq); + if (reqs.size() == group.getGroup().size()) { + return reqs; + } else { + return Collections.emptyList(); + } + } + + private enum SdfGroup { + Sdf01And57(SdfTable.SDF_01, SdfTable.SDF_57), + Sdf08And21(SdfTable.SDF_08, SdfTable.SDF_21); + private final Collection group; + + SdfGroup(SdfTable... group) { + this.group = List.of(group); + } + + public Collection getGroup() { + return group; + } + + public static Optional groupByTable(SdfTable t) { + return Arrays.stream(SdfGroup.values()) + .filter(g -> g.getGroup().contains(t)) + .findFirst(); + } + } +}