SDF06, SDF57 DM*I

This commit is contained in:
ialbert 2023-08-23 11:01:30 +03:00
parent c4476c89a9
commit 0ad02d573b
7 changed files with 121 additions and 19 deletions

View file

@ -265,6 +265,7 @@ public class ValidationConfig {
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
Sdf06NewValidationRule.Fields,
Sdf06NewValidationRule.CompanySymbolAndCompanyPresent,
Sdf06NewValidationRule.AccountPresent,
Sdf06NewValidationRule.Session,

View file

@ -77,7 +77,7 @@ public class AnltSearcher {
if (company == null) {
return new AnltSearch(ClearingError.CompanyNotFoundB, tcr.getCompanyId());
}
return new AnltSearch(infoAcc, company);
return new AnltSearch(infoAcc, company, tcr);
}
private static String getTkrCodeFromComment(String comment) {
@ -97,18 +97,20 @@ public class AnltSearcher {
private EnumMessage error;
private Account account;
private Company company;
private TradingClearingRegistry tcr;
public AnltSearch(IEnumId errSubject, Object... args) {
this.error = new EnumMessage(errSubject, args);
}
public AnltSearch(Account account, Company company) {
public AnltSearch(Account account, Company company, TradingClearingRegistry tcr) {
this.account = account;
this.company = company;
this.tcr = tcr;
}
public boolean isFound() {
return error == null && account != null && company != null;
return error == null && account != null && company != null && tcr != null;
}
public EnumMessage getError() {
@ -134,5 +136,13 @@ public class AnltSearcher {
public void setCompany(Company company) {
this.company = company;
}
public TradingClearingRegistry getTcr() {
return tcr;
}
public void setTcr(TradingClearingRegistry tcr) {
this.tcr = tcr;
}
}
}

View file

@ -22,6 +22,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRe
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
import ru.spcex.clearing.service.registry.DmiService;
import ru.spcex.clearing.service.schedule.TradingTimeService;
import ru.spcex.clearing.service.validation.Sdf06NewValidationRule;
@ -35,7 +36,6 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
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.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidator;
import ru.spcex.platform.utils.validation.ValidatorImpl;
@ -49,6 +49,8 @@ import java.util.Map;
import java.util.Optional;
import java.util.function.Function;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@Service
public class Sdf06Executor {
private final Logger log = LoggerFactory.getLogger(getClass());
@ -57,6 +59,7 @@ public class Sdf06Executor {
private final ImdgProvider imdgProvider;
private final Imdg<Account> accImdg;
private final Imdg<TradingClearingRegistry> tcrImdg;
private final Imdg<Company> cmpImdg;
private final Imdg<SDf06> sdf06Imdg;
private final Imdg<SDf07> sdf07Imdg;
private final ImdgId idGenerator;
@ -95,6 +98,7 @@ public class Sdf06Executor {
this.kafkaSender = kafkaSender;
this.accImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
this.tcrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
this.cmpImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.sdf06Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf06, SDf06.class);
this.sdf07Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf07, SDf07.class);
this.filenameObtainer = filenameObtainer;
@ -154,14 +158,15 @@ public class Sdf06Executor {
statementImdg.insert(stmt);
log.debug("Statement created: {}", stmt.getId());
if (InOutDirection.out.getKey().equals(stmt.getInOutDirection()) && tradingTimeService.isTradingTime()) {
requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode()));
requests.add(GatewayRequestCreator.gatewayRequestPart(stmt, company.getTradingCode(), tcr.getCode()));
} else {
processedApproved(stmt, sDf06, now, sdf07GroupId);
dmiService.setProc(stmt.getAccountId(),
company.getId(),
tcr.getId(),
CurrencyCode.RUB.getKey(),
stmt.getAmount());
safeBD(sDf06.getSum()),
sDf06.getNumber());
sdf07WasCreated = true;
}
}
@ -260,7 +265,8 @@ public class Sdf06Executor {
stmt.getAddresseeId(),
searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null),
CurrencyCode.RUB.getKey(),
stmt.getAmount()
safeBD(sdf06.getSum()),
sdf06.getNumber()
);
} else {
log.trace("statement.id={}, sdf07.id={}, sdf06.id={} rejected (by gateway answer)",
@ -318,13 +324,14 @@ public class Sdf06Executor {
private Statement createStatementBySdf06(SDf06 sdf06, Long companyId, Account account) {
Statement stmt = new Statement();
stmt.setContract(sdf06.getNumber().toString());
stmt.setAddresseeId(companyId);
stmt.setSenderId(Sender.Prc.getId());
stmt.setStatementType(StatementType.incr.getKey());
stmt.setComment(sdf06.getSpec());
stmt.setAccountId(account.getId());
stmt.setAccount(account.getAccount());
BigDecimal amount = BigDecimalUtil.safeBD(sdf06.getSum());
BigDecimal amount = safeBD(sdf06.getSum());
if (amount.compareTo(BigDecimal.ZERO) >= 0) {
stmt.setInOutDirection(InOutDirection.in.getKey());
} else {
@ -354,6 +361,7 @@ public class Sdf06Executor {
ImdgValidationContext<SDf06> ctx = new ImdgValidationContext<>();
ctx.addImdg(IMDGDistributedNames.Map_Account, accImdg);
ctx.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, tcrImdg);
ctx.addImdg(IMDGDistributedNames.Map_Company, cmpImdg);
ctx.setValidatedObject(sdf06);
ValidatorImpl<ImdgValidationContext<SDf06>> v = new ValidatorImpl<>(ctx,
Sdf06NewValidationRule.AccountPresent, Sdf06NewValidationRule.TcrPresent);

View file

@ -18,12 +18,18 @@ 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.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationListRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.AnltSearcher;
import ru.spcex.clearing.service.LoggingService;
import ru.spcex.clearing.service.builder.RegistryBuilder;
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
import ru.spcex.clearing.service.model.Result;
import ru.spcex.clearing.service.registry.DmiService;
import ru.spcex.clearing.service.schedule.TradingTimeService;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.session.stage.util.RegistryUtil;
import ru.spcex.platform.classes.base.SpcexObjectBase;
@ -44,8 +50,10 @@ import java.time.Instant;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.regex.Pattern;
@ -73,12 +81,16 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
private final Imdg<Currency> currencyImdg;
private final AnltSearcher anltSearcher;
private final IMessageResolver messageResolver;
private final DmiService dmiService;
private final TradingTimeService timeService;
private final KafkaSender kafka;
private final Pattern pattern = Pattern.compile("№.*");
public Sdf57Executor(@Qualifier("sdf57Validator") Function<SDf57, IValidator> sDf57Validator,
LoggingService errorLogger,
ImdgProvider imdgProvider,
IMessageResolver errorResolver, AnltSearcher anltSearcher, IMessageResolver messageResolver) {
IMessageResolver errorResolver, AnltSearcher anltSearcher, IMessageResolver messageResolver, DmiService dmiService, TradingTimeService timeService,
@Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafka) {
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.sdf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);
@ -95,6 +107,9 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
this.anltSearcher = anltSearcher;
this.messageResolver = messageResolver;
this.securityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Security, Security.class);
this.dmiService = dmiService;
this.timeService = timeService;
this.kafka = kafka;
}
//todo доделать контроль sdf01 и sdf57
@ -148,6 +163,10 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
if (company == null || account == null) {
return Optional.empty();
}
if (AccountType.Tran.equalsByKey(account.getAccountType())) {
log.trace("sdf57.id={} account type {}, skipping statement creation", sdf57.getId(), AccountType.Tran.getKey());
return Optional.empty();
}
Statement statement = create(sdf57, company, account, inOut);
statementImdg.insert(statement);
return Optional.of(statement);
@ -181,6 +200,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
.companyId(company.getId()) //раньше было stmt.getAddresseeId()
.accountId(account.getId()) //раньше было stmt.getAccountId()
.find();
AtomicReference<Registry> amtCached = new AtomicReference<>();
amtFound.ifPresentOrElse(rgs -> {
log.debug("stmt.id={}, found AM*T.id={}, updating...", stmt.getId(), rgs.getId());
updateReg(stmt, rgs);
@ -209,6 +229,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
rgs.getBalance(),
rgs.getDebit(),
rgs.getCredit());
amtCached.set(rgs);
}, () -> {
Registry registry = RegistryBuilder.builder(imdgProvider)
.statement(stmt)
@ -228,8 +249,10 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
registry.getId(),
registryF.getId(),
registryB.getId());
amtCached.set(registry);
});
InOutDirection direction = IEnumKey.getEnumByKey(InOutDirection.class, stmt.getInOutDirection());
boolean dm_tWasCreated = false;
if (!StringUtils.isEmpty(sdf57.getSpecif()) && pattern.matcher(sdf57.getSpecif()).find() && InOutDirection.in.equals(direction)) {
//todo доделать 11. Идентификация средств УК на клиринговом счете
//stmt.id=340270002, sdf57.specif=Возврат депозита по договору DT1000S001U/280623/8/7 new DM*T.id=339610032
@ -275,6 +298,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
} else if (om_tFound.isPresent()) {
Registry om_t = om_tFound.get();
Registry registryD = copyRegD(om_t, stmt, RegistryUnit.T);
dm_tWasCreated = true;
registryImdg.insert(registryD);
log.debug("stmt.id={}, sdf57.specif={}, found OM*T.id={} new DM*T.id={}",
stmt.getId(),
@ -295,6 +319,14 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
stmt.setOperationStatus(OperationStatus.Executed.getKey());
statementImdg.update(stmt);
log.debug("stmt.id={} status -> {}", stmt.getId(), OperationStatus.Executed.getKey());
if (!dm_tWasCreated && timeService.isTradingTime() && InOutDirection.in.equalsByKey(stmt.getInOutDirection())) {
log.trace("sending gateway request for stmt.id={}", stmt.getId());
kafka.sendRequestToQueue(Consts.ASSET_OPERATION,
gatewayRequest(stmt,
company.getTradingCode(),
amtCached.get().getTradingClearingRegistry()));
}
} else {
log.debug("statement.id={} activness validation failed {}", stmt.getId(), messageResolver.resolve(err.get()));
stmt.setErrorCodeId(err.get().getSubject().getId()); // fixme ErrorText insert
@ -316,6 +348,10 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
searchResult.getAccount().getId(),
searchResult.getCompany().getId());
createRegistryIfNeeded.accept(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount()));
dmiService.setOk(searchResult.getAccount().getId(),
searchResult.getCompany().getId(),
CurrencyCode.RUB.getKey(),
sdf57.getDbfId().toString());
} else {
log.debug("stmt.id={} comment='{}' error: {}. Operating through DMAU registry",
stmt.getId(),
@ -346,6 +382,11 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
});
//DMAU подразумевает увеличение остатка поля balance и credit registry.code=DMAU
}
} else {
dmiService.setOk(acc.getId(),
cmp.getId(),
CurrencyCode.RUB.getKey(),
sdf57.getDbfId().toString());
}
};
statementDeb.map(stmt -> new StmtCmpAcc(stmt, companyDeb, accountDeb)).ifPresent(registersUpdate);
@ -625,4 +666,10 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
R apply(T1 arg1, T2 arg2, T3 arg3);
}
private AssetOperationListRequest gatewayRequest(Statement stmt, String tradingCode, String tcrCode) {
AssetOperationListRequest gatewayRequest = new AssetOperationListRequest();
AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart(stmt, tradingCode, tcrCode);
gatewayRequest.setAssetOperationRequests(List.of(req));
return gatewayRequest;
}
}

View file

@ -0,0 +1,18 @@
package ru.spcex.clearing.service.integration;
import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
import ru.spcex.platform.enumeration.CurrencyCode;
public class GatewayRequestCreator {
public static AssetOperationRequest gatewayRequestPart(Statement stmt, String tradingCode, String tcrCode) {
AssetOperationRequest req = new AssetOperationRequest();
req.setStatementId(stmt.getId());
req.setAmount(stmt.getAmount());
req.setSecuritySymbol(CurrencyCode.RUB.getKey()); //fixme retrieve security symbol from validator
req.setTradingCode(tradingCode);
req.setCode(tcrCode);
req.setDirection(stmt.getInOutDirection());
return req;
}
}

View file

@ -35,35 +35,43 @@ public class DmiService {
Long companyId,
Long tcrId,
String securitySymbol,
BigDecimal balance) {
findDmi(companyId, accountId, securitySymbol).ifPresentOrElse(d__i -> {
BigDecimal balance,
BigDecimal contract) {
findDmi(companyId, accountId, securitySymbol, contract.toString()).ifPresentOrElse(d__i -> {
d__i.setRegistryStatus(RegistryStatus.PROC.getKey());
d__i.setUpdated(Instant.now());
log.trace("updating D**I status to {}", d__i.getRegistryStatus());
rgsImdg.update(d__i);
}, () -> createDmi(companyId, tcrId, securitySymbol, balance, "outDocument"));
}, () -> createDmi(companyId, tcrId, securitySymbol, balance, contract.toString()));
}
public void setOk(Long accountId, Long companyId, Long tcrId, String securitySymbol) {
findDmi(companyId, accountId, securitySymbol).ifPresent(d__i -> {
public void setOk(Long accountId,
Long companyId,
String securitySymbol,
String contract) {
findDmi(companyId, accountId, securitySymbol, contract).ifPresentOrElse(d__i -> {
d__i.setRegistryStatus(RegistryStatus.OK.getKey());
d__i.setUpdated(Instant.now());
log.trace("updating {} status to {}", d__i.getRegistryCode(), d__i.getRegistryStatus());
rgsImdg.update(d__i);
}, () -> {
log.trace("didn't find DM*I by companyId={} accountId={} securitySymbol={} contract={}",
companyId, accountId, securitySymbol, contract);
});
}
private Optional<Registry> findDmi(Long companyId, Long accountId, String securitySymbol) {
private Optional<Registry> findDmi(Long companyId, Long accountId, String securitySymbol, String contract) {
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
RegistryTradingParams rgsCde = RegistryTradingParams.D__I;
ImdgPredicate prdct = pb.and(
pb.sql(RegistryCodeSqlBuilder.getInstance(rgsCde).build()),
pb.equals("companyId", companyId),
pb.equals("accountId", accountId),
pb.equals("securitySymbol", securitySymbol)
pb.equals("securitySymbol", securitySymbol),
pb.equals("contract", contract)
);
log.trace("searching {} by {}", rgsCde, prdct);
log.trace("searching D**I by {}", prdct);
Registry d__i = rgsImdg.getFirstObjectByPredicate(prdct);
return Optional.ofNullable(d__i);
}
@ -79,8 +87,8 @@ public class DmiService {
securitySymbol,
rgsCde);
if (am_t.isEmpty()) {
log.error("{} not found for tcr.id={} company.id={} security.symbol={}",
rgsCde, tcrId, companyId, securitySymbol);
log.error("AM*T not found for tcr.id={} company.id={} security.symbol={}",
tcrId, companyId, securitySymbol);
return;// Optional.empty();
}
Registry rgsD = am_t.get().clone();

View file

@ -31,6 +31,16 @@ import java.util.Optional;
import java.util.function.Function;
public enum Sdf06NewValidationRule implements IValidationRule<ImdgValidationContext<SDf06>> {
Fields() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf06> context) {
SDf06 validatedObject = context.getValidatedObject();
if (validatedObject.getNumber() == null) {
return of(ClearingError.WrongField, "number");
}
return empty();
}
},
CompanySymbolAndCompanyPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf06> context) {