Merge branch 'sdf08' into dev

This commit is contained in:
etreschenkov 2023-05-31 10:45:44 +03:00
commit 79c3c9bdad
6 changed files with 369 additions and 4 deletions

View file

@ -12,6 +12,7 @@ import ru.clearing.classes.statics.data.misc.Session;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.classes.statics.data.sdf.SDf01;
import ru.clearing.classes.statics.data.sdf.SDf08;
import ru.clearing.classes.statics.data.sdf.SDf57;
import ru.clearing.classes.statics.data.security.MoneyMarketSecurity;
import ru.clearing.classes.statics.data.security.Security;
@ -171,4 +172,18 @@ public class ValidationConfig {
);
};
}
@Bean("sdf08ValidatorNew")
public Function<SDf08, IValidator> sdf08ValidatorNew() {
return sDf08 -> {
ImdgValidationContext<SDf08> context = new ImdgValidationContext<>();
context.setValidatedObject(sDf08);
context.addImdg(IMDGDistributedNames.Map_Account, imdgAccount);
context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
return new ValidatorImpl<>(context,
Sdf08NewValidationRule.CompanyPresent,
Sdf08NewValidationRule.AccountPresent
);
};
}
}

View file

@ -56,6 +56,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
this.sdfImdgs.put(SdfTable.SDF_57, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf57, SDf57.class));
this.sdfImdgs.put(SdfTable.SDF_04, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class));
this.sdfImdgs.put(SdfTable.SDF_13, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf13, SDf13.class));
this.sdfImdgs.put(SdfTable.SDF_08, imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class));
}
@Override
@ -86,13 +87,13 @@ public class StatementService extends QueueConsumer implements InitializingBean
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
pairOfSdfRequest.remove(key);
}
} else if (List.of(SdfTable.SDF_08, SdfTable.SDF_13).contains(table)) {
} else if (List.of(SdfTable.SDF_04, SdfTable.SDF_13).contains(table)) {
//not implemented part; it's actually stage number 8 from any session
{
//всегда сначала обработаем sdf04
Long key = completePairKey.get();
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
processSdf08(pair.getFirst());
processSdf04(pair.getFirst());
//затем sdf13
processSdf13(pair.getSecond());
@ -101,6 +102,13 @@ public class StatementService extends QueueConsumer implements InitializingBean
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_SECOND_PART, continueSessionBn);
pairOfSdfRequest.remove(key);
}
} else if (List.of(SdfTable.SDF_08).contains(table)){
{
Long key = completePairKey.get();
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
processSdf08(pair.getFirst());
pairOfSdfRequest.remove(key);
}
}
}
}
@ -113,6 +121,16 @@ public class StatementService extends QueueConsumer implements InitializingBean
service.execute(sdfGroup, statementRequest);
}
private void processSdf04(StatementRequest statementRequest) {
Imdg<SDf04> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class);
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", statementRequest.getGroupId()));
AbstractExecutor service = executorsMap.get(SdfTable.SDF_04);
if (service != null) {
service.execute(sdfGroup, statementRequest);
}
}
private void processSdf08(StatementRequest statementRequest) {
Imdg<SDf08> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf08, SDf08.class);
Collection<? extends SpcexObjectBase> sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
@ -156,6 +174,12 @@ public class StatementService extends QueueConsumer implements InitializingBean
private Optional<Long> saveRequest(StatementRequest statementRequest) {
SdfTable sdfTable = statementRequest.getTable();
//immediately process
if (List.of(SdfTable.SDF_08).contains(sdfTable)){
Long id = imdgProvider.getImdgIdGenerator().nextId();
pairOfSdfRequest.put(id, new Pair<>(statementRequest, null));
return Optional.of(id);
}
if (pairOfSdfRequest.isEmpty()) {
Long id = imdgProvider.getImdgIdGenerator().nextId();
pairOfSdfRequest.put(id, getPairByTableName(sdfTable, statementRequest));

View file

@ -0,0 +1,269 @@
package ru.spcex.clearing.service.executors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClearingAccount;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.classes.statics.data.sdf.SDf08;
import ru.clearing.classes.statics.data.sdf.SDf09;
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.error.ClearingErrorInternal;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.model.Result;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.session.stage.util.RegistryUtil;
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.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
@Service
public class Sdf08Executor extends AbstractExecutor<SDf08> {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Statement> statementImdg;
private final ImdgProvider imdgProvider;
private final Imdg<Registry> registryImdg;
private final Imdg<Security> securityImdg;
private final Function<SDf08, IValidator> sDf08Validator;
private final Imdg<ClearingAccount> clearingAccountImdg;
private final Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
private final Imdg<SDf09> sdf09Imdg;
private final IMessageResolver messageResolver;
public Sdf08Executor(@Qualifier("sdf08ValidatorNew") Function<SDf08, IValidator> sDf08Validator,
ImdgProvider imdgProvider,
IMessageResolver messageResolver) {
this.imdgProvider = imdgProvider;
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class);
this.clearingAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class);
this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
this.sdf09Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf09, SDf09.class);
this.sDf08Validator = sDf08Validator;
this.messageResolver = messageResolver;
}
@Override
public String exportTableName() {
return "";
}
@Override
public boolean isNeedToSendCommand() {
return false;
}
@Override
public void sendCommand(KafkaSender kafkaSender, Result result) {
}
public Result execute(Collection<SDf08> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
for (SDf08 sdf08 : sdf) {
IValidator validator = sDf08Validator.apply(sdf08);
Optional<EnumMessage> error = validator.tillFirstError();
Company company = validator.getStored(ValidationStored.Sdf08Company);
Account account = validator.getStored(ValidationStored.Sdf08Account);
if (statementRequest.getAccountCreationResults().size() == 0
&& ClearingErrorInternal.AccountNotPresent.equals(error.map(EnumMessage::getSubject).orElse(null))) {
result.getAccountRequests().add(createAccountRequestPart(sdf08.getId(), sdf08.getDepoCode(), company.getId()));
log.info("account {} for sdf08.id={} not found - send request for creation", sdf08.getDepoCode(), sdf08.getId());
continue;
} else if (ClearingErrorInternal.AccountNotPresent.equals(error.map(EnumMessage::getSubject).orElse(null))) {
log.error("fatal error: resumed processing after generating accounts, but no account found for sdf01.id={}", sdf08.getId());
}
if (error.isPresent()) {
log.error("sdf01.id={} error: {}", sdf08.getId(), messageResolver.resolve(error.get()));
//fixme инициировать = команда для другого сервиса? sdf02Imdg.insert(createErrorSdf02(sdf01, error.get(), generationIdForGroup));
// sdf02Imdg.insert(createErrorSdf02(sdf01, error.get(), generationIdForGroup));
continue;
}
Statement stmt = createSdf08Statement(sdf08, company, account);
SDf09 sdf09New = createSuccessSdf09(sdf08, generationIdForGroup); //fixme тоже мб убрать
sdf09Imdg.insert(sdf09New);
stmt.setOutSDfId(sdf09New.getId());
statementImdg.insert(stmt);
Optional<EnumMessage> err = validateActiveness(company, account);
if (err.isEmpty()) {
Optional<Registry> reg = findReg(stmt);
Registry rgs;
if (reg.isPresent()) {
rgs = reg.get();
updateReg(stmt, rgs);
registryImdg.update(rgs);
} else {
rgs = createRegistryByStatement(stmt, company, account);
registryImdg.insert(rgs);
}
//todo
//2. отправить notification на backend
stmt.setOperationStatus(OperationStatus.Executed.getKey());
statementImdg.update(stmt);
} else {
stmt.setErrorCodeId(err.get().getSubject().getId()); // fixme ErrorText insert
stmt.setOperationStatus(OperationStatus.Rejected.getKey());
statementImdg.update(stmt);
}
}
return result;
}
private Optional<EnumMessage> validateActiveness(Company company, Account account) {
if (!WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
return Optional.of(new EnumMessage(ClearingError.CompanyNotActive, company.getId()));
}
if (!WorkflowStatus.Active.equalsByKey(account.getStatus())) {
return Optional.of(new EnumMessage(ClearingError.AccountNotActive, account.getId()));
}
return Optional.empty();
}
private AccountSdfRequestPart createAccountRequestPart(Long sdf01Id, String account, Long companyId) {
AccountSdfRequestPart req = new AccountSdfRequestPart();
req.setAccount(account);
req.setCompanyId(companyId);
req.setAccountType(AccountType.Clrn.getKey());
req.setSdfId(sdf01Id);
return req;
}
private Statement createSdf08Statement(SDf08 sdf08, Company company, Account account) {
Statement stmt = new Statement();
stmt.setAddresseeId(company.getId());
stmt.setAddresseeId(company.getId());
stmt.setSenderId(Sender.Rdc.getId());
stmt.setStatementType(StatementType.incr.getKey());
stmt.setAccountId(account.getId());
stmt.setAccount(account.getAccount());
stmt.setInOutDirection(InOutDirection.in.getKey());
if (sdf08.getQuantity() != null) {
stmt.setAmount(new BigDecimal(sdf08.getQuantity()));
}
stmt.setOperationStatus(OperationStatus.Pending.getKey());
stmt.setInSDfId(sdf08.getId());
stmt.setInOutSDfType(InOutSDfType.type1.getKey());
stmt.setClearingDate(LocalDate.now());
Security security = securityImdg.getSingleObjectByFieldValues(Map.of("securitySymbol", sdf08.getSecurityCode()));
if (security != null) {
stmt.setSecurityId(security.getId());
}
stmt.setCreated(Instant.now());
return stmt;
}
private SDf09 createSuccessSdf09(SDf08 sdf08, Long generationIdForGroup) {
SDf09 sDf09 = new SDf09();
sDf09.setInDocument(sdf08.getId().toString());
sDf09.setDepoCode(sdf08.getDepoCode());
sDf09.setQuantity(sdf08.getQuantity());
sDf09.setSecurityCode(sdf08.getSecurityCode());
sDf09.setClientName(sdf08.getClientName());
sDf09.setGenerationId(generationIdForGroup);
sDf09.setGenerationTime(Instant.now());
sDf09.setResult("OK!");
return sDf09;
}
private Registry createRegistryByStatement(Statement statement, Company company, Account account) {
Registry rgs = new Registry();
rgs.setCompanyId(statement.getAddresseeId());
rgs.setTradingCode(company.getTradingCode());
rgs.setClearingCode(company.getClearingCode());
rgs.setShortName(company.getShortName());
rgs.setFullName(company.getFullName());
rgs.setAccountId(account.getId());
rgs.setAccountType(account.getAccountType());
rgs.setAccount(account.getAccount());
rgs.setRegistryDesignation(RegistryDesignation.A.getKey());
rgs.setRegistryInstrumentType(RegistryInstrumentType.S.getKey());
ClearingAccount accountForStatement = clearingAccountImdg.getSingleObjectByID(statement.getAccountId());
if (accountForStatement != null) {
rgs.setRegistryCapacity(accountForStatement.getClearingAccountType());
}
rgs.setRegistryUnit(RegistryUnit.T.getKey());
rgs.setRegistryCode(RegistryUtil.clearingCode(rgs));
Collection<TradingClearingRegistry> tcrsByAccount = tradingClearingRegistryImdg.getCollectionObjectsByFieldValues(Map.of(
"moneyAccountId", statement.getAccountId(),
"companyId", statement.getAddresseeId()
));
if (!tcrsByAccount.isEmpty()) {
TradingClearingRegistry tcr = tcrsByAccount.iterator().next();
rgs.setTradingClearingRegistryId(tcr.getId());
rgs.setTradingClearingRegistry(tcr.getCode());
}
rgs.setRegistryStatus(RegistryStatus.PROC.getKey());
rgs.setSecurityId(statement.getSecurityId());
if (statement.getSecurityId() != null) {
Security security = securityImdg.getSingleObjectByID(statement.getSecurityId());
if (security != null) {
rgs.setSecuritySymbol(security.getSecuritySymbol());
}
}
rgs.setCheckBalance(BigDecimalUtil.safeBD(statement.getAmount()));
rgs.setDiffBalance(rgs.getCheckBalance().negate());
rgs.setBalanceDimension(BalanceDimension.MONY.getKey()); //fixme ! смотри описание и ссылка на начало html'ки
//fixme !rgs.setSettlementCode();
rgs.setTradingDate(statement.getSettlementDate()); //fixme ! today ?
rgs.setClearingDate(LocalDate.now());
//fixme rgs.setRefundDate();
//fixme rgs.setValueDate();
//создается на базе stmt, companyCred, accountDeb
rgs.setCounterPartyId(statement.getAddresseeId());
rgs.setCreated(Instant.now());
return rgs;
}
private void updateReg(Statement s, Registry r) {
r.setCheckBalance(s.getAmount());
BigDecimal balance = BigDecimalUtil.safeBD(r.getBalance());
BigDecimal checkBalance = BigDecimalUtil.safeBD(r.getCheckBalance());
r.setDiffBalance(balance.subtract(checkBalance));
r.setUpdated(Instant.now());
}
private Optional<Registry> findReg(Statement s) {
RegistryTradingParams p = new RegistryTradingParams(
RegistryDesignation.A, RegistryInstrumentType.S, null, RegistryUnit.T
);
String sql = RegistryCodeSqlBuilder.getInstance(p).build();
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql),
pb.sql(sql),
pb.equals("accountId", s.getAccountId()),
pb.equals("companyId", s.getAddresseeId())
);
return Optional.ofNullable(registryImdg.getSingleObjectByPredicate(rgstrPredicate));
}
}

View file

@ -0,0 +1,54 @@
package ru.spcex.clearing.service.validation;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.sdf.SDf08;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.error.ClearingErrorInternal;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Optional;
public enum Sdf08NewValidationRule implements IValidationRule<ImdgValidationContext<SDf08>> {
CompanyPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf08> context) {
SDf08 sdf08 = context.getValidatedObject();
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
if (TextUtil.isEmpty(sdf08.getClientName())) {
of(ClearingError.CompanyNotFoundB, sdf08.getClientName());
}
Company company = companyImdg.getSingleObjectBySQL("fullName = '" + sdf08.getClientName() + "'");
if (company == null) {
return of(ClearingError.CompanyNotFoundB, sdf08.getClientName());
}
context.storeObject(ValidationStored.Sdf08Company, company);
return empty();
}
}, AccountPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf08> context) {
SDf08 sdf08 = context.getValidatedObject();
Imdg<Account> accountImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class);
Account acc = accountImdg.getSingleObjectBySQL("account = '" + sdf08.getDepoCode()
+ "' and accountType='" + ru.spcex.platform.enumeration.AccountType.Depo.getKey() + "'");
if (acc == null) {
return of (ClearingErrorInternal.AccountNotPresent, sdf08.getDepoCode());
}
context.storeObject(ValidationStored.Sdf08Account, acc);
return empty();
}
};
@Override
public String ruleName() {
return "Sdf01NewValidationRule." + name();
}
}

View file

@ -6,5 +6,7 @@ public enum ValidationStored {
Sdf57CompanyDeb, Sdf57CompanyCred, Sdf57AccountDeb, Sdf57AccountCred,
Sdf01Company, Sdf01Account
Sdf01Company, Sdf01Account,
Sdf08Company, Sdf08Account
}

View file

@ -4,7 +4,8 @@ import ru.spcex.platform.utils.enumeration.IEnumId;
public enum Sender implements IEnumId {
One(1L), /* СПВБ */
Prc(2L); /* ПРЦ */
Prc(2L), /* ПРЦ */
Rdc(3L); /* РДЦ */
private final Long key;