etreschenkov 2023-05-31 10:30:07 +03:00
parent 14dc721cfd
commit bfc487cfa5
2 changed files with 160 additions and 30 deletions

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

@ -5,11 +5,15 @@ 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;
@ -17,13 +21,16 @@ 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;
@ -38,18 +45,26 @@ import java.util.function.Function;
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;
}
@ -93,39 +108,45 @@ public class Sdf08Executor extends AbstractExecutor<SDf08> {
}
Statement stmt = createSdf08Statement(sdf08, company, account);
//на данном шаге company существует -> getId ok
//формируем пакетный запрос на добавление account
//ответ придет в этот же метод, process
// result.getAccountRequests().add(createAccountRequestPart(sdf01.getId(), sdf01.getAccount(), company.getId()));
// log.info("account {} for sdf01.id={} not found - send request for creation", sdf01.getAccount(), sdf01.getId());
// continue;
// log.debug("Process sdf08 record; sdf08.id: {}", sdf08.getId());
// Collection<Registry> registries = selectRegistryForSDF04(sdf08.getC_acc_cred());
// registries.forEach(registry -> unlockRegistry(registry, new BigDecimal(sdf08.getPay_val())));
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;
}
protected Collection<Registry> selectRegistryForSDF04(String account) {
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate query = pb.and(
pb.and(
pb.equals("registryDesignation", RegistryDesignation.A.getKey()),
pb.equals("registryInstrumentType", RegistryInstrumentType.M.getKey()),
pb.equals("registryUnit", RegistryUnit.B.getKey())
),
pb.equals("account", account)
);
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(query);
log.trace("Selected {} registry's by sql: {}", result.size(), query);
return result;
}
boolean unlockRegistry(Registry registry, BigDecimal value) {
if (registry.getBalance() == null) registry.setBalance(BigDecimal.ZERO);
registry.setBalance(registry.getBalance().subtract(value));
return true;
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) {
@ -160,4 +181,89 @@ public class Sdf08Executor extends AbstractExecutor<SDf08> {
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));
}
}