From 6fdfcb041fb441e45db6556f0a82e3153e7e2163 Mon Sep 17 00:00:00 2001 From: etreschenkov Date: Wed, 17 May 2023 12:05:26 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-290 stage 4 --- .../clearing/session/stage/TaskType.java | 4 + .../stage/impl/InclusionObligations.java | 80 +++++++++++++++++++ .../stage/task/InclusionToPoolPayload.java | 22 +++++ 3 files changed, 106 insertions(+) create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java create mode 100644 clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/InclusionToPoolPayload.java diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java index 86b6b740d..c1fb4a584 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java @@ -11,5 +11,9 @@ public enum TaskType { * step 1 */ DealsPrepare, + /** + * step 4 + */ + InclusionToPool, } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java new file mode 100644 index 000000000..4a61cf6c9 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java @@ -0,0 +1,80 @@ +package ru.spcex.clearing.session.stage.impl; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.session.stage.ISessionStage; +import ru.spcex.clearing.session.stage.StageResult; +import ru.spcex.clearing.session.stage.Task; +import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.time.LocalDate; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + + +@Service +public class InclusionObligations implements ISessionStage { + private final Logger log = LoggerFactory.getLogger(getClass()); + //todo remove (set all in single method setImdg(provider -> setImdg1();setIdGenerator();...) + private ImdgProvider imdgProvider; + private ImdgId idGenerator; + private Imdg registryImdg; + private KafkaSender kafkaSender; + + @Autowired + public InclusionObligations(ImdgProvider imdgProvider) { + this.imdgProvider = imdgProvider; + this.idGenerator = imdgProvider.getImdgIdGenerator(); + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + } + + @Override + public StageResult submit(Task task) { + InclusionToPoolPayload payload = (InclusionToPoolPayload) task.getData(); + switch (task.getTaskType()) { + case InclusionToPool -> { + return inclusionToPool(payload.getSessionType(), payload.getCounterPartyId()); + } + default -> { + throw new IllegalStateException("Unknown task type: " + task.getTaskType()); + } + } + } + + + private StageResult inclusionToPool(String sessionType, Long counterPartyId) { + String sqlCondition = String.format("registryCode like '%s' and " + + "registryStatus = '%s' and " + + "settlementDate = '%s' and " + + "sessionType = '%s'", + "[O/T][S/M][*][T]", + "PROC", + LocalDate.now(), + sessionType); + + Collection obligations = registryImdg.getCollectionObjectsBySQL(sqlCondition); + Map> registryByGroupId = obligations.stream(). + filter(registry -> registry.getCompanyId().equals(counterPartyId)). + collect(Collectors.groupingBy(Registry::getGroupId)); + + for (Map.Entry> entrySet : registryByGroupId.entrySet()) { + log.debug("Processing set of registry with groupId: {}", entrySet.getKey()); + for (Registry registry : entrySet.getValue()){ + registry.setRegistryStatus("POOL"); + registryImdg.update(registry); + } + } + + return new StageResult(null, true); + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/InclusionToPoolPayload.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/InclusionToPoolPayload.java new file mode 100644 index 000000000..41bba14c6 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/InclusionToPoolPayload.java @@ -0,0 +1,22 @@ +package ru.spcex.clearing.session.stage.task; + +public class InclusionToPoolPayload { + private String sessionType; + private Long counterPartyId; + + public String getSessionType() { + return sessionType; + } + + public void setSessionType(String sessionType) { + this.sessionType = sessionType; + } + + public Long getCounterPartyId() { + return counterPartyId; + } + + public void setCounterPartyId(Long counterPartyId) { + this.counterPartyId = counterPartyId; + } +}