ReturnDepositSession

This commit is contained in:
ialbert 2023-06-01 15:44:39 +03:00
parent c040188512
commit 8955aafdb9
3 changed files with 443 additions and 0 deletions

View file

@ -0,0 +1,206 @@
package ru.spcex.clearing.session.stage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.session.stage.impl.*;
import ru.spcex.clearing.session.stage.task.EndStageNotificationPayload;
import ru.spcex.clearing.session.stage.task.FinishingSessionPayload;
import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload;
import ru.spcex.clearing.session.stage.task.InspectionPoolPayload;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.SessionStatus;
import ru.spcex.platform.enumeration.SessionType;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import java.time.LocalDate;
import java.util.Collection;
@Service
public class ReturnDepositSession extends AbstractSession implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final BalanceRevise balanceRevise;
private final DealsPrepare dealsPrepare;
private final RequirementsAndObligationCreation requirementsAndObligationCreation;
private final ObligationAdmission obligationsAdmission;
private final InclusionObligations inclusionObligations;
private final InspectionObligationsDepositReturn inspectionObligationsDepositReturn;
private final FormingRegistersOnOS formingRegistersOnOS;
private final FormingPaymentInstruction formingPaymentInstruction;
private final UnlockResources unlockResources;
private final FinishingSession finishingSession;
private final EndStageNotification endStageNotification;
private final Imdg<Registry> registryImdg;
public ReturnDepositSession(
ImdgProvider imdgProvider,
BalanceRevise balanceRevise,
DealsPrepare dealsPrepare,
RequirementsAndObligationCreation requirementsAndObligationCreation,
ObligationAdmission obligationsAdmission,
InclusionObligations inclusionObligations,
FormingRegistersOnOS formingRegistersOnOS,
FormingPaymentInstruction formingPaymentInstruction,
UnlockResources unlockResources,
FinishingSession finishingSession,
EndStageNotification endStageNotification,
IMessageResolver messageResolver,
InspectionObligationsDepositReturn inspectionObligationsDepositReturn) {
super(imdgProvider, messageResolver);
this.balanceRevise = balanceRevise;
this.dealsPrepare = dealsPrepare;
this.requirementsAndObligationCreation = requirementsAndObligationCreation;
this.obligationsAdmission = obligationsAdmission;
this.inclusionObligations = inclusionObligations;
this.formingRegistersOnOS = formingRegistersOnOS;
this.formingPaymentInstruction = formingPaymentInstruction;
this.unlockResources = unlockResources;
this.finishingSession = finishingSession;
this.endStageNotification = endStageNotification;
this.inspectionObligationsDepositReturn = inspectionObligationsDepositReturn;
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
}
@Override
public void afterPropertiesSet() throws Exception {
ImdgPredicateBuilder rgsPrctBuilder = registryImdg.predicateBuilder();
inclusionObligations.addRegistryCondition(
rgsPrctBuilder.or(rgsPrctBuilder.equals("registryStatus", RegistryStatus.PROC.getKey()),
rgsPrctBuilder.equals("registryStatus", RegistryStatus.MNG.getKey()))
);
inclusionObligations.addRegistryCondition(
rgsPrctBuilder.less("valueDate", LocalDate.now())
);
}
@Override
public void runSession(BaseRequest<?> req) {
if (!startSession()) {
return;
}
StageResult<?> submit = balanceRevise.submit(new Task<>(TaskType.StartRevise, null));
if (!submit.success) {
log.error("stage {} error={}", balanceRevise.getClass().getSimpleName(), messageResolver.resolve(submit.error));
endSession();
} else {
log.info("stage BalanceRevise success, waiting for a response from kafka");
}
}
public void continueSession(BaseRequest<?> req) {
try {
if (!isRunning()) {
return;
}
if (checkStage(TaskType.StartRevise)) {
firstPart(req);
} else if (checkStage(TaskType.FormingPaymentInstruction)) {
finishPart(req);
} else {
log.info("will not continue session, current stage is {}", currStage.get());
}
} catch (StageException e) {
//already logged
}
}
private void firstPart(BaseRequest<?> req) {
try {
//stage 4
{
InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload();
inclusionToPoolPayload.setSessionType(currSession.getSessionType());
runStage(TaskType.InclusionToPool, inclusionToPoolPayload, inclusionObligations);
}
//stage 5
{
InspectionPoolPayload companyIdPayload = new InspectionPoolPayload();
companyIdPayload.setProcessedCompanyId(currSession.getCompanyId());
runStage(TaskType.InspectionObligations, companyIdPayload, inspectionObligationsDepositReturn);
}
//stage 6
runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection<Registry>
//stage 7
StageResult<Collection<PaymentInstruction>> paymentResult = runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction);
if (paymentResult.getStageResult().isEmpty()) {
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
// finishPart(req);
}
} catch (StageException e) {
//already logged
}
}
public void finishPart(BaseRequest<?> req) {
try {
if (!checkStage(TaskType.FormingPaymentInstruction)) {
log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException();
}
//stage 10
{
FinishingSessionPayload payload = new FinishingSessionPayload();
payload.setSessionId(currSession.getId());
runStage(TaskType.FinishingSession, payload, finishingSession);
}
//stage 11
{
EndStageNotificationPayload payload = new EndStageNotificationPayload();
payload.setSection(currSession.getSection());
payload.setSessionId(currSession.getId());
runStage(TaskType.EndStageNotification, payload, endStageNotification);
}
} catch (StageException e) {
//already logged
}
}
private boolean startSession() {
synchronized (this.currStage) {
if (this.currStage.get() != null) {
log.info("already running session.id={}", this.currSession.getId());
return false;
} else {
TaskType startStatus = TaskType.StartRevise;
Session newSession = new Session();
newSession.setSection(section().getKey());
newSession.setSessionType(sectionType().getKey());
newSession.setSessionStatus(startStatus.getKey());
newSession.setWorkflowStatus(SessionStatus.ACTV.getKey());
newSession.setClearingDate(LocalDate.now());
//todo companyId/securityId/userId передается из сообщения очереди
sessionImdg.insert(newSession);
currSession = newSession;
log.info("started new session.id={}", this.currSession.getId());
currStage.set(TaskType.StartRevise);
return true;
}
}
}
@Override
protected Section section() {
return Section.MKR;
}
@Override
protected SessionType sectionType() {
return SessionType.XDEP;
}
}

View file

@ -0,0 +1,236 @@
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.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service
public class InspectionObligationsDepositReturn implements ISessionStage {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Registry> registryImdg;
private final static RegistryTradingParams OS_T;
private final static RegistryTradingParams OM_T;
private final static RegistryTradingParams TS_T;
private final static RegistryTradingParams TM_T;
private final static RegistryTradingParams DM_X;
private final static RegistryTradingParams DM_T;
private final static RegistryTradingParams AM_F;
static {
OS_T = new RegistryTradingParams(RegistryDesignation.O,
RegistryInstrumentType.S,
null,
RegistryUnit.T);
OM_T = new RegistryTradingParams(RegistryDesignation.O,
RegistryInstrumentType.M,
null,
RegistryUnit.T);
TS_T = new RegistryTradingParams(RegistryDesignation.T,
RegistryInstrumentType.S,
null,
RegistryUnit.T);
TM_T = new RegistryTradingParams(RegistryDesignation.T,
RegistryInstrumentType.M,
null,
RegistryUnit.T);
DM_X = new RegistryTradingParams(RegistryDesignation.D,
RegistryInstrumentType.M,
null,
RegistryUnit.X);
DM_T = new RegistryTradingParams(RegistryDesignation.D,
RegistryInstrumentType.M,
null,
RegistryUnit.T);
AM_F = new RegistryTradingParams(RegistryDesignation.A,
RegistryInstrumentType.M,
null,
RegistryUnit.F);
}
@Autowired
public InspectionObligationsDepositReturn(ImdgProvider imdgProvider) {
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
}
@Override
public StageResult<?> submit(Task<?> task) {
switch (task.getTaskType()) {
case InspectionObligations -> {
return inspectionObligations();
}
default -> throw new IllegalStateException("Unknown task type: " + task.getTaskType());
}
}
private StageResult<?> inspectionObligations() {
String sqlCondition = String.format("(%s) and registryStatus = '%s'",
RegistryCodeSqlBuilder.getInstance(OS_T, OM_T, TS_T, TM_T).build(),
RegistryStatus.POOL.getKey());
Map<Long, List<Registry>> registriesByGroup = registryImdg.getCollectionObjectsBySQL(sqlCondition)
.stream().
collect(Collectors.groupingBy(Registry::getGroupId));
log.info("registry groups found {}", registriesByGroup.size());
for (Map.Entry<Long, List<Registry>> entry : registriesByGroup.entrySet()) {
List<Registry> group = entry.getValue();
Optional<Registry> omtInGroupO = group.stream().filter(registry -> equalByRgs(OM_T, registry)).findFirst();
if (omtInGroupO.isEmpty()) {
log.error("groupId {} failed to find OM*T registry", entry.getKey());
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG));
continue;
}
Registry omtRgs = omtInGroupO.get();
Optional<Registry> amfAssetO = searchAssetByOMT(omtRgs);
if (amfAssetO.isEmpty()) {
log.error("OM*T register.id={} groupId={} failed to find AM*F asset", omtRgs.getId(), entry.getKey());
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG));
continue;
}
Registry amfAsset = amfAssetO.get();
log.debug("groupId={}, OM*T.id={}, AM*F.id={}", entry.getKey(), omtRgs.getId(), amfAsset.getId());
BigDecimal omtBalance = safeBD(omtRgs.getBalance());
BigDecimal amfBalance = safeBD(amfAsset.getBalance());
Optional<Registry> dmx = searchDmx(omtRgs);
if (dmx.isPresent()) {
BigDecimal dmxBalance = safeBD(dmx.get().getBalance());
log.debug("OM*T.id={} -> DM*X.id={}, dmxBalance={}, omtBalance={}, amfBalance={}",
omtRgs.getId(),
dmx.get().getId(),
dmxBalance, omtBalance, amfBalance
);
if (dmxBalance.compareTo(omtBalance) >= 0 && amfBalance.compareTo(omtBalance) >= 0) {
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK));
} else {
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG));
}
continue;
}
Optional<Registry> dmtInfo = searchDmtInfo(omtRgs);
if (dmtInfo.isPresent()) {
BigDecimal dmtBalance = safeBD(dmtInfo.get().getBalance());
log.debug("OM*T.id={} -> DM*T(INFO).id={}, dmtBalance={}, omtBalance={}, amfBalance={}",
omtRgs.getId(),
dmtInfo.get().getId(),
dmtBalance, omtBalance, amfBalance
);
if (dmtBalance.compareTo(omtBalance) >= 0 && amfBalance.compareTo(omtBalance) >=0) {
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK));
} else {
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG));
}
continue;
}
Optional<Registry> dmtClnr = searchDmtClrn(omtRgs);
if (dmtClnr.isPresent()) {
BigDecimal dmtBalance = safeBD(dmtClnr.get().getBalance());
log.debug("OM*T.id={} -> DM*T(CLNR).id={}, dmtBalance={}, omtBalance={}, amfBalance={}",
omtRgs.getId(),
dmtClnr.get().getId(),
dmtBalance, omtBalance, amfBalance
);
if (dmtBalance.compareTo(omtBalance) >= 0 && amfBalance.compareTo(omtBalance) >= 0) {
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.OK));
} else {
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG));
}
continue;
}
log.info("OM*T.id={} -> no DM*X/DM*T(INFO/CLRN) registry found. setting MNG status to group", omtRgs.getId());
group.forEach(rgs -> updateStatus(rgs, RegistryStatus.MNG));
}
return new StageResult<>(null, true);
}
private void updateStatus(Registry registry, RegistryStatus registryStatus) {
log.trace("Update registry.id: {} to {}", registry.getId(), registryStatus.getKey());
registry.setRegistryStatus(registryStatus.getKey());
registry.setUpdated(Instant.now());
registryImdg.update(registry);
}
/**
* Если УК зачисляет средства на ТБС Инициатора
*/
private Optional<Registry> searchDmx(Registry rgs) {
String sqlCondition = String.format("(%s) and companyId = %d",
RegistryCodeSqlBuilder.getInstance(DM_X).build(),
rgs.getCounterPartyId());
Registry dmx = registryImdg.getSingleObjectBySQL(sqlCondition);
return Optional.ofNullable(dmx);
}
/**
* Если УК зачисляет средства на свой регистр на КС
*/
private Optional<Registry> searchDmtInfo(Registry rgs) {
String sqlCondition = String.format("(%s) and accountType='%s' and companyId = %d and counterPartyId = %d",
RegistryCodeSqlBuilder.getInstance(DM_T).build(),
AccountType.Info.getKey(),
rgs.getCompanyId(),
rgs.getCounterPartyId());
Registry dmt = registryImdg.getSingleObjectBySQL(sqlCondition);
return Optional.ofNullable(dmt);
}
/**
* Если УК зачисляет средства на свой ТБС
*/
private Optional<Registry> searchDmtClrn(Registry rgs) {
String sqlCondition = String.format("(%s) and accountType='%s' and companyId = %d and counterPartyId = %d",
RegistryCodeSqlBuilder.getInstance(DM_T).build(),
AccountType.Clrn.getKey(),
rgs.getCompanyId(),
rgs.getCounterPartyId());
Registry dmt = registryImdg.getSingleObjectBySQL(sqlCondition);
return Optional.ofNullable(dmt);
}
private Optional<Registry> searchAssetByOMT(Registry obligation) {
String sqlCondition = String.format("%s and " +
"tradingClearingRegistryId = '%s' and " +
"companyId = '%s'",
RegistryCodeSqlBuilder.getInstance(AM_F).build(),
obligation.getTradingClearingRegistryId(),
obligation.getCompanyId());
Registry amf = registryImdg.getSingleObjectBySQL(sqlCondition);
return Optional.ofNullable(amf);
}
private static boolean equalByRgs(RegistryTradingParams rgsParams, Registry rgs) {
return rgsParams.equalByRegistry(
IEnumKey.getEnumByKey(RegistryDesignation.class, rgs.getRegistryDesignation()),
IEnumKey.getEnumByKey(RegistryInstrumentType.class, rgs.getRegistryInstrumentType()),
IEnumKey.getEnumByKey(RegistryCapacity.class, rgs.getRegistryCapacity()),
IEnumKey.getEnumByKey(RegistryUnit.class, rgs.getRegistryUnit())
);
}
}

View file

@ -8,6 +8,7 @@ public enum RegistryUnit implements IEnumKey {
R("R"),
F("F"),
B("B"),
X("X"),
;
private final String key;