multiple actions

This commit is contained in:
ialbert 2024-10-01 18:34:22 +03:00
parent 4ccca75aad
commit 05d1135ae2
7 changed files with 842 additions and 12 deletions

View file

@ -24,6 +24,8 @@ import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.clearing.session.state.SsnEvent;
import ru.spcex.clearing.session.state.action.CreateSessionAction;
import ru.spcex.clearing.session.state.action.DealsPrepareAction;
import ru.spcex.clearing.session.state.action.FormingPaymentInstructionAssetsAction;
import ru.spcex.clearing.session.state.action.FormingPaymentInstructionReturnMkrAction;
import ru.spcex.clearing.session.state.action.InclusionToPoolAction;
import ru.spcex.clearing.session.state.action.InspectionObligationsDepositReturnAction;
import ru.spcex.clearing.session.state.action.InspectionObligationsV2Action;
@ -55,6 +57,8 @@ public class FinalSessionStateMachineConfig
private final InclusionToPoolAction inclusionToPoolAction;
private final InspectionObligationsDepositReturnAction inspOblDepositReturnAction;
private final InspectionObligationsV2Action inspectionObligationsV2Action;
private final FormingPaymentInstructionReturnMkrAction formingPaymentInstructionReturnMkrAction;
private final FormingPaymentInstructionAssetsAction formingPaymentInstructionAssetsAction;
@Autowired
@ -62,7 +66,16 @@ public class FinalSessionStateMachineConfig
ImdgProvider imdgProvider,
@Qualifier("marketCodesForBn")
Supplier<List<String>> marketCodes,
Sdf56Action sendSdf56Action, ReviseStage1 reviseStage1Action, DealsPrepareAction dealsPrepareAction, RequirementAndObligationCreationAction reqAndOblAction, ObligationAdmissionAction obligationAdmissionAction, InclusionToPoolAction inclusionToPoolAction, InspectionObligationsDepositReturnAction inspOblDepositReturnAction, InspectionObligationsV2Action inspectionObligationsV2Action
Sdf56Action sendSdf56Action,
ReviseStage1 reviseStage1Action,
DealsPrepareAction dealsPrepareAction,
RequirementAndObligationCreationAction reqAndOblAction,
ObligationAdmissionAction obligationAdmissionAction,
InclusionToPoolAction inclusionToPoolAction,
InspectionObligationsDepositReturnAction inspOblDepositReturnAction,
InspectionObligationsV2Action inspectionObligationsV2Action,
FormingPaymentInstructionReturnMkrAction formingPaymentInstructionReturnMkrAction,
FormingPaymentInstructionAssetsAction formingPaymentInstructionAssetsAction
) {
this.sendSdf56Action = sendSdf56Action;
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
@ -74,6 +87,8 @@ public class FinalSessionStateMachineConfig
this.inclusionToPoolAction = inclusionToPoolAction;
this.inspOblDepositReturnAction = inspOblDepositReturnAction;
this.inspectionObligationsV2Action = inspectionObligationsV2Action;
this.formingPaymentInstructionReturnMkrAction = formingPaymentInstructionReturnMkrAction;
this.formingPaymentInstructionAssetsAction = formingPaymentInstructionAssetsAction;
//stages settings:
ImdgPredicateBuilder rgsPrctBuilder = rgsImdg.predicateBuilder();
@ -143,11 +158,13 @@ public class FinalSessionStateMachineConfig
.withExternal()
.source(TaskType.InclusionToPool).target(TaskType.InspectionObligations)
.action(inspOblDepositReturnAction)
.action(inspectionObligationsV2Action)
.and()
//internal (no state change) triggerless(without an event) action
.withInternal()
.source(TaskType.InspectionObligations)
.action(inspectionObligationsV2Action);
.withExternal()
.source(TaskType.InspectionObligations).target(TaskType.FormingPaymentInstruction)
.action(formingPaymentInstructionReturnMkrAction)
.action(formingPaymentInstructionAssetsAction)
;

View file

@ -6,4 +6,7 @@ public enum DataEnum {
session, //Session
dealsPrepared, //List<ExecutionCommon>
counterPartyId, //Long
paymentInstructionReturnMkr, //List<PaymentInstruction>
paymentInfo,
}

View file

@ -0,0 +1,595 @@
package ru.spcex.clearing.session.state.action;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.function.Function;
import java.util.function.Supplier;
import java.util.stream.Stream;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Scope;
import org.springframework.statemachine.StateContext;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.SettlementHouseProperties;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.sdf.SDf03;
import ru.clearing.classes.statics.data.sdf.SDf12;
import ru.clearing.classes.statics.data.security.Security;
import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.Sdf03Creator;
import ru.spcex.clearing.service.builder.PaymentInstructionBuilderV2;
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
import ru.spcex.clearing.service.registry.PaymentStateMarkService;
import ru.spcex.clearing.service.registry.RegistryManager;
import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.clearing.session.stage.impl.PaymentInfo;
import ru.spcex.clearing.session.stage.util.RegistryUtil;
import ru.spcex.clearing.session.state.DataEnum;
import ru.spcex.clearing.session.state.SsnEvent;
import ru.spcex.platform.enumeration.AccountStatus;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.Allowed;
import ru.spcex.platform.enumeration.CurrencyCode;
import ru.spcex.platform.enumeration.InstrumentType;
import ru.spcex.platform.enumeration.OperationStatus;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
import static ru.spcex.platform.enumeration.RegistryTradingParams.*;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.enumeration.SdfTable;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Sender;
import ru.spcex.platform.enumeration.SessionType;
import ru.spcex.platform.enumeration.StatementType;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector;
import ru.spcex.platform.utils.collection.Pair;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.enumeration.SimpleMessageResolver;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class FormingPaymentInstructionAssetsAction extends AbstractSessionActionForOkErrorHandling {
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 Imdg<Statement> stmtImdg;
private Imdg<PaymentInstruction> paymentInstructionImdg;
private Imdg<Security> securityImdg;
private Imdg<Account> accountImdg;
private Imdg<Company> companyImdg;
private Imdg<SDf03> sDf03Imdg;
private Imdg<SDf12> sDf12Imdg;
private Imdg<SettlementHouseProperties> stlmHPropsImdg;
private final Sdf03Creator sdf03Creator;
private KafkaSender kafkaSender;
private final IMessageResolver msgResolver = new SimpleMessageResolver();
private final SecuritySelector<Security> securitySelector;
private final RegistryManager rgsMng;
private Section section;
private SessionType sessionType;
private final PaymentStateMarkService dmvSrv;
private final AssetTBFProcessing assets;
@Autowired
public FormingPaymentInstructionAssetsAction(ImdgProvider imdgProvider,
Sdf03Creator sdf03Creator, KafkaSender kafkaSender, RegistryManager rgsMng, PaymentStateMarkService dmvSrv, AssetTBFProcessing assets) {
this.sdf03Creator = sdf03Creator;
this.kafkaSender = kafkaSender;
this.imdgProvider = imdgProvider;
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class);
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.sDf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class);
this.sDf12Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf12, SDf12.class);
this.paymentInstructionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.stlmHPropsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SettlementHouseProperties, SettlementHouseProperties.class);
this.stmtImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.rgsMng = rgsMng;
this.securitySelector = new SecuritySelector<>(imdgProvider, Security.class);
this.dmvSrv = dmvSrv;
this.assets = assets;
}
@Override
protected void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
if (SessionType.FINL.equals(sessionType) || SessionType.MEDM.equals(sessionType) || SessionType.UNIT.equals(sessionType)) {
dmvSrv.createDmv(sessionId);
}
Collection<Registry> AMBregistries = selectAMBRegistries(sessionId);
log.debug("found AMB registries.size() = {}", AMBregistries.size());
Instant now = Instant.now();
Account dtrnAcc = accountImdg.getFirstObjectBySQL("accountType = '%s' and status = '%s' and processingSign = '%s'"
.formatted(AccountType.Dtrn.getKey(), AccountStatus.ACTIVE.getKey(), Allowed.ALLOWED.getKey()));
if (dtrnAcc == null) {
throw new RuntimeException(
"%s (%d): accountType = %s".formatted(
ClearingError.AccountNotPresent.name(),
ClearingError.AccountNotPresent.getId(),
AccountType.Dtrn.getKey())
);
}
List<Pair<PaymentInstruction, Registry>> paymentInstructions = new ArrayList<>();
for (Registry registry : AMBregistries) {
Account tranAcc = accountImdg.getFirstObjectBySQL(("accountType = '%s' " +
"and status = '%s' " +
//fixme empty for RUB
"and currency = '%s' " +
"and processingSign = '%s'")
.formatted(AccountType.Tran.getKey(),
AccountStatus.ACTIVE.getKey(),
registry.getSecuritySymbol(),
Allowed.ALLOWED.getKey()));
if (tranAcc == null) {
throw new RuntimeException(
"%s (%d): accountType = %s".formatted(
ClearingError.AccountNotPresent.name(),
ClearingError.AccountNotPresent.getId(),
AccountType.Tran.getKey())
);
}
//по каждому AMB регистру создаем PaymentInstruction
//проверяем balance посчитанный на шаге 5
Account counterAcc = null;
if (AccountType.Clrn.equalsByKey(registry.getAccountType())) {
counterAcc = accountImdg.getSingleObjectByID(registry.getAccountId());
} else if (AccountType.Info.equalsByKey(registry.getAccountType())) {
counterAcc = accountImdg.getFirstObjectByFieldValues(
Map.of("companyId", 1L,
"currency", registry.getSecuritySymbol(),
"accountType", AccountType.Anlt.getKey())
);
}
Long senderId;
Long addresseeId;
Account debitLegAccount;
Account creditLegAccount;
BigDecimal amount;
amount = safeBD(registry.getBalance()).add(adjustAmountByReturns(registry, sessionId));
boolean isPositiveBalance = amount.compareTo(BigDecimal.ZERO) > 0;
if (amount.compareTo(BigDecimal.ZERO) == 0) {
continue;
}
if (isPositiveBalance) {
senderId = registry.getCompanyId();
addresseeId = Sender.One.getId();
debitLegAccount = tranAcc;
creditLegAccount = counterAcc;
Optional<Registry> payerAmtO = rgsMng.findRelatedAsset(
registry.getTradingClearingRegistryId(),
registry.getCompanyId(),
registry.getSecuritySymbol(),
AM_T);
payerAmtO.ifPresent(amt -> {
amt.setSettledDebit(safeBD(amt.getSettledDebit()).add(safeBD(registry.getBalance().abs())));
setUpdatedStoreInImdg(amt, now);
});
} else {
senderId = Sender.One.getId();
addresseeId = registry.getCompanyId();
debitLegAccount = counterAcc;
creditLegAccount = tranAcc;
Optional<Registry> payerAmtO = rgsMng.findRelatedAsset(
registry.getTradingClearingRegistryId(),
registry.getCompanyId(),
registry.getSecuritySymbol(),
AM_T
);
payerAmtO.ifPresent(amt -> {
amt.setSettledCredit(safeBD(amt.getSettledCredit()).add(safeBD(registry.getBalance().abs())));
setUpdatedStoreInImdg(amt, now);
});
}
amount = amount.abs();
PaymentInstructionBuilderV2 paymentInstructionBuilder = PaymentInstructionBuilderV2.builder(imdgProvider)
.registry(registry)
.sender(senderId)
.addressee(addresseeId)
.debitLegAccount(debitLegAccount)
.creditLegAccount(creditLegAccount)
.amount(amount)
.currency(registry.getSecuritySymbol())
.sessionId(sessionId)
.checkBLKD((SessionType.FINL.equals(sessionType) || SessionType.MEDM.equals(sessionType))
&& !isPositiveBalance ? registry : null)
.purpose(String.format("Перевод по итогу клиринга по ТКР %s", registry.getTradingClearingRegistry()));
PaymentInstruction paymentInstruction = paymentInstructionBuilder.build();
log.debug("Created paymentInstruction by registry.id: {}", registry.getId());
paymentInstructionImdg.insert(paymentInstruction);
registry.setPaymentId(paymentInstruction.getId());
registry.setUpdated(now);
registryImdg.update(registry);
paymentInstructions.add(new Pair<>(paymentInstruction, registry));
}
if (section == null || !section.equals(Section.MKR)) {
Collection<Registry> ASBregistries = selectASBRegistries(sessionId);
log.debug("found ASB registries.size() = {}", ASBregistries.size());
for (Registry registry : ASBregistries) {
//по каждому AMB регистру создаем PaymentInstruction
//проверяем balance посчитанный на шаге 5
if (registry.getBalance().compareTo(BigDecimal.ZERO) == 0) {
log.debug("Skip creating paymentInstruction by registry with 0 balance");
continue;
}
Account counterAcc = accountImdg.getSingleObjectByID(registry.getAccountId());
boolean isPositiveBalance = registry.getBalance().compareTo(BigDecimal.ZERO) > 0;
Long senderId;
Long addresseeId;
Account debitLegAccount;
Account creditLegAccount;
BigDecimal amount;
if (isPositiveBalance) {
senderId = registry.getCompanyId();
addresseeId = Sender.One.getId();
debitLegAccount = dtrnAcc;
creditLegAccount = counterAcc;
amount = registry.getBalance() == null ? null : registry.getBalance().abs();
Optional<Registry> payerAstO = rgsMng.findRelatedAsset(
registry.getTradingClearingRegistryId(),
registry.getCompanyId(),
registry.getSecuritySymbol(),
AS_T
);
payerAstO.ifPresent(amt -> {
amt.setSettledDebit(safeBD(amt.getSettledDebit()).add(safeBD(registry.getBalance().abs())));
setUpdatedStoreInImdg(amt, now);
});
} else {
senderId = Sender.One.getId();
addresseeId = registry.getCompanyId();
debitLegAccount = counterAcc;
creditLegAccount = dtrnAcc;
amount = registry.getBalance() == null ? null : registry.getBalance().abs();
Optional<Registry> payerAstO = rgsMng.findRelatedAsset(
registry.getTradingClearingRegistryId(),
registry.getCompanyId(),
registry.getSecuritySymbol(),
AS_T
);
payerAstO.ifPresent(amt -> {
amt.setSettledCredit(safeBD(amt.getSettledCredit()).add(safeBD(registry.getBalance().abs())));
setUpdatedStoreInImdg(amt, now);
});
}
PaymentInstructionBuilderV2 paymentInstructionBuilder = PaymentInstructionBuilderV2.builder(imdgProvider)
.registry(registry)
.sender(senderId)
.addressee(addresseeId)
.debitLegAccount(debitLegAccount)
.creditLegAccount(creditLegAccount)
.amount(amount)
.sessionId(sessionId)
.currency(SessionType.CURR.equals(sessionType) || Section.CURR.equalsByKey(registry.getSection())
? registry.getSecuritySymbol() : null)
.purpose(String.format("Перевод по итогу клиринга по ТКР %s", registry.getTradingClearingRegistry()));
PaymentInstruction paymentInstruction = paymentInstructionBuilder.build();
log.debug("Created paymentInstruction by registry.id: {}", registry.getId());
paymentInstructionImdg.insert(paymentInstruction);
registry.setPaymentId(paymentInstruction.getId());
registry.setUpdated(now);
registryImdg.update(registry);
paymentInstructions.add(new Pair<>(paymentInstruction, registry));
}
}
@SuppressWarnings("unchecked")
List<PaymentInstruction> pmtAlreadyCreated = ctx
.getExtendedState()
.get(DataEnum.paymentInstructionReturnMkr, List.class);
pmtAlreadyCreated = pmtAlreadyCreated != null ? pmtAlreadyCreated : Collections.emptyList();
if (!pmtAlreadyCreated.isEmpty()) {
log.debug("PaymentInstructions return size {}, PaymentInstructions deals size {}. sending SDF03",
pmtAlreadyCreated.size(),
paymentInstructions.size());
paymentInstructions = Stream.concat(
pmtAlreadyCreated
.stream()
.map(p -> new Pair<PaymentInstruction, Registry>(p, null)),
paymentInstructions.stream()).toList();
}
Collection<SdfTable> sdfTables = sendSdfs(paymentInstructions, sessionId, dtrnAcc);
PaymentInfo paymentInfo = new PaymentInfo(paymentInstructions.stream().map(Pair::getFirst).toList(), sdfTables);
ctx
.getExtendedState()
.getVariables()
.put(DataEnum.paymentInfo, paymentInfo);
}
private Collection<Registry> selectAMBRegistries(Long sessionId) {
ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder();
RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(
AM_B
);
ImdgPredicate rgsCodePrdct = rgsPb.and(
rgsPb.sql(registryCodeSqlBuilder.build()),
rgsPb.equals("sessionId", sessionId)
);
return registryImdg.getCollectionObjectsByPredicate(rgsCodePrdct);
}
private Collection<Registry> selectASBRegistries(Long sessionId) {
ImdgPredicateBuilder rgsPb = registryImdg.predicateBuilder();
RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(
AS_B
);
ImdgPredicate rgsCodePrdct = rgsPb.and(
rgsPb.sql(registryCodeSqlBuilder.build()),
rgsPb.equals("sessionId", sessionId)
);
return registryImdg.getCollectionObjectsByPredicate(rgsCodePrdct);
}
private BigDecimal adjustAmountByReturns(Registry amb, Long sessionId) {//companyid, tcr/account, sessionId, RUB
if (SessionType.UNIT.equals(sessionType)) return BigDecimal.ZERO;
Long companyId = amb.getCompanyId();
Long tradingClearingRegistryId = amb.getTradingClearingRegistryId();
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate cL = pb.and(
pb.equals("companyId", companyId),
pb.equals("tradingClearingRegistryId", tradingClearingRegistryId),
pb.equals("registryStatus", RegistryStatus.OK.getKey()),
pb.equals("sessionId", sessionId),
pb.or(
pb.sql(RegistryCodeSqlBuilder.getInstance(OM_T).build()),
pb.sql(RegistryCodeSqlBuilder.getInstance(TM_T).build())
)
);
BigDecimal reduce = registryImdg.getCollectionObjectsByPredicate(cL)
.stream()
.filter(rgs -> rgs.getValueDate() != null)
.filter(rgs -> rgs.getSettlementDate() != null)
.filter(rgs -> rgs.getSettlementDate().isAfter(rgs.getValueDate()))
.map(rgs -> {
if (RegistryDesignation.T.equalsByKey(rgs.getRegistryDesignation())) {
return rgs.getBalance();
} else {
return rgs.getBalance().abs().negate();
}
})
.reduce(BigDecimal.ZERO, BigDecimal::add);
log.debug("{}.id={} {} adjust value {}",
amb.getRegistryCode(),
amb.getId(),
cL,
reduce);
return reduce;
}
private void setUpdatedStoreInImdg(Registry rgs, Instant now) {
rgs.setUpdated(now);
registryImdg.update(rgs);
}
private Collection<SdfTable> sendSdfs(List<Pair<PaymentInstruction, Registry>> formedPaymentInstructions, Long sessionId, Account dtrnAcc) {
Collection<SdfTable> sentSdfs = new ArrayList<>();
List<SDf03> sDf03Created = new ArrayList<>();
List<SDf12> sDf12Created = new ArrayList<>();
for (Pair<PaymentInstruction, Registry> pair : formedPaymentInstructions) {
PaymentInstruction paymentInstruction = pair.getFirst();
Registry rgs = pair.getSecond();
Account account = accountImdg.getSingleObjectByID(paymentInstruction.getCreditLeg_accountId());
Security security = securityImdg.getSingleObjectByID(paymentInstruction.getCreditLeg_securityId());
if (SessionType.CURR.equals(sessionType) || (rgs != null
&& RegistryInstrumentType.M.equalsByKey(rgs.getRegistryInstrumentType())
&& !CurrencyCode.isRub(rgs.getSecuritySymbol()))) {
boolean madeStatements = statementsWereMade(paymentInstruction, rgs);
if (madeStatements) {
continue;
}
}
if (List.of(AccountType.Corr, AccountType.Clrn, AccountType.Tran, AccountType.Anlt, AccountType.Info)
.contains(IEnumKey.getEnumByKey(AccountType.class, account.getAccountType()))) {
//если info подставить anlt (единственный счет в системе)
sDf03Created.add(sdf03Creator.create(paymentInstruction));
} else if (!InstrumentType.CRNC.equalsByKey(security.getInstrumentType()) &&
List.of(AccountType.Depo, AccountType.Dtrn).contains(IEnumKey.getEnumByKey(AccountType.class, account.getAccountType()))) {
sDf12Created.add(newSDf12(paymentInstruction, dtrnAcc));
}
}
Long sdf03GroupId = !sDf03Created.isEmpty() ? imdgProvider.getImdgIdGenerator().nextId() : null;
for (SDf03 sDf03 : sDf03Created) {
sDf03.setGenerationId(sdf03GroupId);
sDf03Imdg.insert(sDf03);
}
if (sdf03GroupId != null) {
ExportToFileRequest exportToFileRequest = new ExportToFileRequest();
exportToFileRequest.setNameOfTable("DF-03");
exportToFileRequest.setSdfGroupId(sdf03GroupId);
kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportToFileRequest);
}
Long sdf12GroupId = null;
Long maxTxNumber = 1L;
if (!sDf12Created.isEmpty()) {
sdf12GroupId = imdgProvider.getImdgIdGenerator().nextId();
ImdgPredicateBuilder predicateBuilder = sDf12Imdg.predicateBuilder();
ImdgPredicate notEmptyTransactionNum = predicateBuilder.not(predicateBuilder.equals("transactionNumber", ""));
Long maxId = sDf12Imdg.aggregateLongMax("id", notEmptyTransactionNum);
if (maxId != null) {
SDf12 sDf12 = sDf12Imdg.getSingleObjectByID(maxId);
maxTxNumber = Long.parseLong(sDf12.getTransactionNumber()) + 1;
}
}
for (SDf12 sDf12 : sDf12Created) {
sDf12.setGenerationId(sdf12GroupId);
sDf12.setTransactionNumber(maxTxNumber.toString());
sDf12.setTransactionQuantity(String.valueOf(sDf12Created.size()));
sDf12Imdg.insert(sDf12);
}
if (sdf12GroupId != null) {
SwtExporterRequest swtExporterRequest = new SwtExporterRequest();
swtExporterRequest.setType("SDF_12");
swtExporterRequest.setGroupId(sdf12GroupId);
swtExporterRequest.setSessionId(sessionId);
kafkaSender.sendRequestToQueue(Consts.SWT_EXPORTER, swtExporterRequest);
}
if (sdf03GroupId != null) sentSdfs.add(SdfTable.SDF_03);
if (sdf12GroupId != null) sentSdfs.add(SdfTable.SDF_12);
return sentSdfs;
}
private SDf12 newSDf12(PaymentInstruction paymentInstruction, Account dtrnAcc) {
log.debug("creating sdf12");
SDf12 sDf12 = new SDf12();
if (Objects.equals(paymentInstruction.getDebitLeg_accountId(), dtrnAcc.getId())) {
sDf12.setDirection("RECFREE");
sDf12.setDepoCodeSender(paymentInstruction.getDebitLeg_account());
sDf12.setDepoCodeAdressee(paymentInstruction.getCreditLeg_account());
} else if (Objects.equals(paymentInstruction.getCreditLeg_accountId(), dtrnAcc.getId())) {
sDf12.setDirection("DELFREE");
sDf12.setDepoCodeSender(paymentInstruction.getCreditLeg_account());
sDf12.setDepoCodeAdressee(paymentInstruction.getDebitLeg_account());
}
sDf12.setId(idGenerator.nextId());
sDf12.setOutDocument(sDf12.getId().toString());
sDf12.setQuantity(BigDecimalUtil.limitDecimalPlaces(paymentInstruction.getCreditLeg_amount().toString(), 2));
Security security = securitySelector.selectSecurityById(paymentInstruction.getCreditLeg_securityId());
if (security != null) {
sDf12.setSecurityCode(security.getSecuritySymbol());
}
// sDf12.setTransactionNumber();
sDf12.setGenerationTime(Instant.now());
return sDf12;
}
private boolean statementsWereMade(PaymentInstruction pmt, Registry am_b) {
ImdgPredicateBuilder pb = stlmHPropsImdg.predicateBuilder();
String curCode = pmt.getCreditLeg_currencyCode();
{
SettlementHouseProperties sttHs = stlmHPropsImdg.getFirstObjectByPredicate(
pb.and(
pb.in("currencyCode", curCode),
pb.equals("companyId", Sender.Prc.getId())
)
);
if (sttHs != null) {
return false;
}
}
Supplier<Statement> stmtBldr = () -> {
Statement stmt = new Statement();
stmt.setSenderId(Sender.One.getId());
stmt.setStatementType(StatementType.incr.getKey());
stmt.setComment("Расчеты по сессии");
stmt.setSettlementDate(LocalDate.now());
stmt.setOperationStatus(OperationStatus.Executed.getKey());//todo тут что-то с апдейтом статуса на EXEC когда создались регистры с нулевым диффом?
stmt.setClearingDate(LocalDate.now());
stmt.setCreated(Instant.now());
stmt.setComment(pmt.getPaymentPurpose());
return stmt;
};
Statement stmtDeb = stmtBldr.get();
Statement stmtCred = stmtBldr.get();
stmtCred.setAddresseeId(pmt.getSenderId());
stmtDeb.setAddresseeId(pmt.getAddresseeId());
stmtDeb.setAccountId(pmt.getDebitLeg_accountId());
stmtDeb.setAccount(pmt.getDebitLeg_account());
stmtCred.setAccountId(pmt.getCreditLeg_accountId());
stmtCred.setAccount(pmt.getCreditLeg_account());
stmtDeb.setSecurityId(pmt.getDebitLeg_securityId());
stmtCred.setSecurityId(pmt.getCreditLeg_securityId());
stmtDeb.setInOutDirection(pmt.getDebitLeg_direction());
stmtCred.setInOutDirection(pmt.getCreditLeg_direction());
stmtDeb.setAmount(pmt.getDebitLeg_amount());
stmtCred.setAmount(pmt.getCreditLeg_amount());
stmtImdg.insert(stmtDeb);
stmtImdg.insert(stmtCred);
if (am_b != null) {
Function<RegistryUnit, Registry> copy = (unit) -> {
Registry res = am_b.clone();
RegistryManager.zeroState(res);
res.setRegistryUnit(RegistryUnit.B.getKey());
res.setRegistryCode(RegistryUtil.clearingCode(res));
return res;
};
Registry am_t = assets.searchByTcrCompanyAccount(am_b, AM_T).orElseThrow(() -> {
log.error("{}.id={} PaymentInstruction.id={} no AM*T found", am_b.getRegistryCode(), am_b.getId(), pmt.getId());
return new RuntimeException("couldn't find AM*T");
});
BigDecimal am_bBalance = safeBD(am_b.getBalance());
am_b.setBalance(BigDecimal.ZERO);
{
BigDecimal amtBalance = safeBD(am_t.getBalance());
am_t.setBalance(amtBalance.subtract(am_bBalance));
}
Registry am_f = assets.searchByTcrCompanyAccount(am_b, AM_F)
.orElseGet(() -> copy
.andThen(r -> {r.setBalance(am_t.getBalance()); return r;})
.apply(RegistryUnit.F)
);
assets.process(am_b, am_t, am_f, BigDecimal.ZERO);
registryImdg.update(am_t);
}
log.debug(("paymentInstruction.id=%d created statements [credit id=%d] and [debit id=%d]. " +
"set am_b.id=%d balance to zero").formatted(
pmt.getId(), stmtCred.getId(), stmtDeb.getId(), am_b != null ? am_b.getId() : null
));
return true;
}
public void setSection(Section section) {
this.section = section;
}
public void setSessionType(SessionType sessionType) {
this.sessionType = sessionType;
}
}

View file

@ -0,0 +1,202 @@
package ru.spcex.clearing.session.state.action;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.context.annotation.Scope;
import org.springframework.statemachine.StateContext;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.builder.PaymentInstructionBuilderFinalMkrDeals;
import ru.spcex.clearing.service.registry.RegistryManager;
import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.clearing.session.state.DataEnum;
import ru.spcex.clearing.session.state.SsnEvent;
import ru.spcex.platform.enumeration.AccountStatus;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.enumeration.Allowed;
import ru.spcex.platform.enumeration.RegistryCapacity;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import static ru.spcex.platform.enumeration.RegistryTradingParams.AM_B;
import static ru.spcex.platform.enumeration.RegistryTradingParams.AM_T;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.collection.Pair;
import ru.spcex.platform.utils.enumeration.IEnumKey;
@Service
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class FormingPaymentInstructionReturnMkrAction extends AbstractSessionActionForOkErrorHandling {
private final Logger log = LoggerFactory.getLogger(getClass());
//todo remove (set all in single method setImdg(provider -> setImdg1();setIdGenerator();...)
private final ImdgProvider imdgProvider;
private final Imdg<Registry> registryImdg;
private final Imdg<PaymentInstruction> paymentInstructionImdg;
private final Imdg<Account> accountImdg;
private final RegistryManager rgsMng;
@Autowired
public FormingPaymentInstructionReturnMkrAction(ImdgProvider imdgProvider,
RegistryManager rgsMng) {
this.imdgProvider = imdgProvider;
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.paymentInstructionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
this.rgsMng = rgsMng;
}
@Override
protected void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
Instant now = Instant.now();
RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.OM_T, RegistryTradingParams.TM_T);
String registryCodeCondition = registryCodeSqlBuilder.build();
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.sql(registryCodeCondition),
pb.equals("registryStatus", RegistryStatus.OK.getKey()),
pb.equals("sessionId", sessionId)
);
Collection<Registry> rgsAll = registryImdg.getCollectionObjectsByPredicate(prdct)
.stream()
.filter(rgs -> rgs.getValueDate() != null)
.filter(rgs -> rgs.getSettlementDate() != null)
.filter(rgs -> rgs.getSettlementDate().isAfter(rgs.getValueDate()))
.toList();
log.debug("LM*T and CM*T size = {}", rgsAll.size());
log.debug("changing assets by CM*T");
List<Registry> requirementsByMoney = rgsAll.stream()
.filter(registry -> equalByRgs(RegistryTradingParams.TM_T, registry))
.toList();
//изменяем активы по требованиям по деньгам CM_T
for (Registry requirementByMoney : requirementsByMoney) {
String sql = String.format("tradingClearingRegistryId = %s and companyId = %s",
requirementByMoney.getTradingClearingRegistryId(), requirementByMoney.getCompanyId());
Collection<Registry> relatedRegistries = registryImdg.getCollectionObjectsBySQL(sql);
log.trace("found {} related registries (by tcrId&companyId) for CM*T.id={}", relatedRegistries.size(), requirementByMoney.getId());
for (Registry relatedRegistry : relatedRegistries) {
if (equalByRgs(AM_T, relatedRegistry)) {
relatedRegistry.setSettledDebit(safeBD(relatedRegistry.getSettledDebit()).add(safeBD(requirementByMoney.getBalance())));
log.trace("updated related registry.id={} for CM*T.id={}", relatedRegistry.getId(), requirementByMoney.getId());
registryImdg.update(relatedRegistry);
}
}
}
List<PaymentInstruction> pmtCreated = new ArrayList<>();
Map<Long, List<Registry>> registriesByGroup = rgsAll
.stream().
collect(Collectors.groupingBy(Registry::getGroupId));
log.info("registry groups found {}", registriesByGroup.size());
for (Map.Entry<Long, List<Registry>> entry : registriesByGroup.entrySet()) {
List<Registry> groupRgs = entry.getValue();
Optional<Registry> lmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.OM_T, rgs)).findFirst();
Optional<Registry> cmtO = groupRgs.stream().filter(rgs -> equalByRgs(RegistryTradingParams.TM_T, rgs)).findFirst();
if (lmtO.isEmpty() || cmtO.isEmpty()) {
log.error("LM*T or CM*T not found for group {}", entry.getKey());
continue;
}
Registry lm_t = lmtO.get(); //obligation by money
Registry cm_t = cmtO.get();
Optional<Registry> dmx = rgsMng.searchDmxByCounterPartyNotOk(lm_t);
if (dmx.isPresent()) {
log.debug("LM*T#id={}, DM*X#id={} found, no action needed for group {}, skipping liability",
lm_t.getId(), dmx.get().getId(), lm_t.getGroupId());
dmx.get().setRegistryStatus(RegistryStatus.OK.getKey());
setUpdatedStoreInImdg(dmx.get(), now);
continue;
}
Optional<Registry> dmtClnr = rgsMng.searchDmtClrnNotOk(lm_t);
Account tranAcc = accountImdg.getFirstObjectBySQL(("accountType = '%s'" +
" and status = '%s'" +
" and processingSign = '%s'" +
" and currency = '%s'")
.formatted(AccountType.Tran.getKey(),
AccountStatus.ACTIVE.getKey(),
Allowed.ALLOWED.getKey(),
lm_t.getSecuritySymbol()));
if (tranAcc == null) {
throw new RuntimeException(
"%s (%d): accountType = %s".formatted(
ClearingError.AccountNotPresent.name(),
ClearingError.AccountNotPresent.getId(),
AccountType.Tran.getKey())
);
}
{
//изменение активов - блокируем средства беред отправкой sdf'ов
Optional<Registry> receiverAmtO = rgsMng.findRelatedAsset(cm_t.getTradingClearingRegistryId(), cm_t.getCompanyId(), lm_t.getSecuritySymbol(), AM_B);
log.debug("changing A* registers based on LM_T.id={} and CM_T.id={} found AM*T(receiver).id={}",
lm_t.getId(),
cm_t.getId(),
receiverAmtO.map(Registry::getId).orElse(null)
);
receiverAmtO.ifPresent(amt -> {
// у отправителя и получателя одинаково, см. в FormingPaymentInstruction
amt.setSettledDebit(safeBD(amt.getSettledDebit()).add(safeBD(lm_t.getBalance())));
setUpdatedStoreInImdg(amt, now);
});
}
Pair<PaymentInstruction, PaymentInstruction> pmts = PaymentInstructionBuilderFinalMkrDeals.builder(imdgProvider)
.lm_t(lm_t)
.cm_t(cm_t)
.tranAcc(tranAcc)
.sessionId(sessionId)
.paymentPurpose("Возврат депозита " + cm_t.getContract() + " по ТКР " + cm_t.getTradingClearingRegistry() + " по итогу клиринга")
.paymentPurposeLmt("Возврат депозита " + lm_t.getContract() + " по ТКР " + lm_t.getTradingClearingRegistry() + " по итогу клиринга")
.currency(lm_t.getSecuritySymbol())
.build();
Pair.forEach(pmts, paymentInstructionImdg::insert);
Pair.forEach(pmts, pmtCreated::add);
log.debug("LM*T#id={}, CM*T#id={} found, PaymentInstruction id={} and id={} created",
lm_t.getId(), cm_t.getId(), pmts.getFirst().getId(), pmts.getSecond().getId());
if (dmtClnr.isPresent()) {
dmtClnr.get().setRegistryStatus(RegistryStatus.OK.getKey());
setUpdatedStoreInImdg(dmtClnr.get(), now);
}
}
ctx.getExtendedState().getVariables().put(DataEnum.paymentInstructionReturnMkr, pmtCreated);
}
private void setUpdatedStoreInImdg(Registry rgs, Instant now) {
rgs.setUpdated(now);
registryImdg.update(rgs);
}
private BigDecimal safeBD(BigDecimal value) {
return value != null ? value : BigDecimal.ZERO;
}
private static boolean equalByRgs(RegistryTradingParams code, Registry rgs) {
return code.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

@ -2,4 +2,5 @@ package ru.spcex.clearing.session.teststate.config;
public enum Event {
startSession, continueRevise,
triggerInternal
}

View file

@ -4,5 +4,6 @@ public enum State {
initial,
startRevise,
continueRevise,
newState
stateWithInternalTransition,
stateAfterInternalTransition,
}

View file

@ -40,8 +40,9 @@ public class TestStateMachineConfig extends EnumStateMachineConfigurerAdapter<St
.initial(State.initial, StateMachineUtil.chain(new InitialAction(), new SelfAction()))
.state(State.initial)
.state(State.startRevise)
.state(State.continueRevise)
.state(State.newState)
.state(State.continueRevise, Event.triggerInternal)
.state(State.stateWithInternalTransition)
.state(State.stateAfterInternalTransition)
;
}
@ -69,16 +70,26 @@ public class TestStateMachineConfig extends EnumStateMachineConfigurerAdapter<St
.and()
.withExternal()
.source(State.continueRevise)
.target(State.newState)
.target(State.stateWithInternalTransition)
.action(context -> {
log.info("action without an event");
log.info("ACTION without an event");
// context.getStateMachine().sendEvent(Event.triggerInternal);
})
.and()
.withInternal()
.source(State.newState)
.source(State.stateWithInternalTransition)
.action(context -> {
log.info("INTERNAL ACTION WITHOUT AN EVENT");
log.info("ACTION INTERNAL without an event");
})
.action(context -> {
log.info("ACTION INTERNAL second without an event");
})
.and()
.withExternal()
.source(State.stateWithInternalTransition)
.timerOnce(1)
.target(State.stateAfterInternalTransition)
.action(ctx -> log.info("ACTION after internal transition"))
;
}