This commit is contained in:
etreschenkov 2023-05-17 12:05:26 +03:00
parent fe79868712
commit 6fdfcb041f
3 changed files with 106 additions and 0 deletions

View file

@ -11,5 +11,9 @@ public enum TaskType {
* step 1
*/
DealsPrepare,
/**
* step 4
*/
InclusionToPool,
}

View file

@ -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<Registry> 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<Registry> obligations = registryImdg.getCollectionObjectsBySQL(sqlCondition);
Map<Long, List<Registry>> registryByGroupId = obligations.stream().
filter(registry -> registry.getCompanyId().equals(counterPartyId)).
collect(Collectors.groupingBy(Registry::getGroupId));
for (Map.Entry<Long, List<Registry>> 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);
}
}

View file

@ -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;
}
}