SDF54/SDF55/SDF57 - D**I регистры
This commit is contained in:
parent
b19912f430
commit
93ffe36ea6
16 changed files with 247 additions and 91 deletions
|
|
@ -44,7 +44,7 @@ public class KafkaConfig {
|
||||||
gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs());
|
gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs());
|
||||||
gatewayConsumer.setAutoOffsetReset(consumer.getAutoOffsetReset());
|
gatewayConsumer.setAutoOffsetReset(consumer.getAutoOffsetReset());
|
||||||
gatewayConsumer.setEnableAutoCommit(consumer.getEnableAutoCommit());
|
gatewayConsumer.setEnableAutoCommit(consumer.getEnableAutoCommit());
|
||||||
gatewayConsumer.setGroupId(consumer.getGroupId() + "-session-" + groupId.getAndIncrement());
|
gatewayConsumer.setGroupId(consumer.getGroupId() + "-gateway-" + groupId.getAndIncrement());
|
||||||
return KafkaConsumerFactory.consumer(gatewayConsumer);
|
return KafkaConsumerFactory.consumer(gatewayConsumer);
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -398,6 +398,7 @@ public class ValidationConfig {
|
||||||
ctx.setValidatedObject(pmtOut);
|
ctx.setValidatedObject(pmtOut);
|
||||||
ctx.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
|
ctx.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
|
||||||
ctx.addImdg(IMDGDistributedNames.Map_Account, imdgAccount);
|
ctx.addImdg(IMDGDistributedNames.Map_Account, imdgAccount);
|
||||||
|
ctx.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
|
||||||
return new ValidatorImpl<>(ctx,
|
return new ValidatorImpl<>(ctx,
|
||||||
PaymentOutboundValidationRule.RequiredFields,
|
PaymentOutboundValidationRule.RequiredFields,
|
||||||
IdPresentRule.instance("senderId",
|
IdPresentRule.instance("senderId",
|
||||||
|
|
@ -413,7 +414,9 @@ public class ValidationConfig {
|
||||||
ClearingError.RequiredFieldEmpty,
|
ClearingError.RequiredFieldEmpty,
|
||||||
ClearingError.DictionaryNotFound),
|
ClearingError.DictionaryNotFound),
|
||||||
PaymentOutboundValidationRule.CreditLegAccount,
|
PaymentOutboundValidationRule.CreditLegAccount,
|
||||||
PaymentOutboundValidationRule.DebitLegAccount);
|
PaymentOutboundValidationRule.DebitLegAccount,
|
||||||
|
PaymentOutboundValidationRule.AddresseePresent,
|
||||||
|
PaymentOutboundValidationRule.TcrPresent);
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeRefundDateRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeRefundDateRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeStatusExtractRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryReturnDepositRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryReturnDepositRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistrySplitDepositActionRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistrySplitDepositActionRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||||
|
|
@ -155,6 +156,9 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
callback(RegistryReturnDepositRequest.class)
|
callback(RegistryReturnDepositRequest.class)
|
||||||
.setFunction(registryService::returnDeposit)
|
.setFunction(registryService::returnDeposit)
|
||||||
.forDestination(Consts.REGISTRY_RETURN_DEPOSIT_ACTION, callbacks::put);
|
.forDestination(Consts.REGISTRY_RETURN_DEPOSIT_ACTION, callbacks::put);
|
||||||
|
callback(RegistryChangeStatusExtractRequest.class)
|
||||||
|
.setFunction(registryService::changeStatusExtract)
|
||||||
|
.forDestination(Consts.REGISTRY_CHANGE_STATUS_EXTRACT_ACTION, callbacks::put);
|
||||||
callback(RegistryChangeRefundDateRequest.class)
|
callback(RegistryChangeRefundDateRequest.class)
|
||||||
.setFunction(registryService::changeRefundDate)
|
.setFunction(registryService::changeRefundDate)
|
||||||
.forDestination(Consts.REGISTRY_CHANGE_REFUND_DATE_ACTION, callbacks::put);
|
.forDestination(Consts.REGISTRY_CHANGE_REFUND_DATE_ACTION, callbacks::put);
|
||||||
|
|
|
||||||
|
|
@ -152,11 +152,11 @@ public class Sdf06Executor {
|
||||||
statementImdg.insert(stmt);
|
statementImdg.insert(stmt);
|
||||||
log.debug("Statement created: {}", stmt.getId());
|
log.debug("Statement created: {}", stmt.getId());
|
||||||
if (InOutDirection.out.getKey().equals(stmt.getInOutDirection()) && tradingTimeService.isTradingTime()) {
|
if (InOutDirection.out.getKey().equals(stmt.getInOutDirection()) && tradingTimeService.isTradingTime()) {
|
||||||
requests.add(GatewayRequestCreator.gatewayRequestPart(stmt, company.getTradingCode(), tcr.getCode()));
|
requests.add(GatewayRequestCreator.from(stmt, company.getTradingCode(), tcr.getCode()));
|
||||||
} else {
|
} else {
|
||||||
processedApproved(stmt, sDf06, now, sdf07GroupId);
|
processedApproved(stmt, sDf06, now, sdf07GroupId);
|
||||||
if (tcr != null) {
|
if (tcr != null) {
|
||||||
dmiService.setProc(tcr.getId(),
|
dmiService.setProcContract(tcr.getId(),
|
||||||
CurrencyCode.RUB.getKey(),
|
CurrencyCode.RUB.getKey(),
|
||||||
safeBD(sDf06.getSum()),
|
safeBD(sDf06.getSum()),
|
||||||
sDf06.getNumber().toString());
|
sDf06.getNumber().toString());
|
||||||
|
|
@ -263,7 +263,7 @@ public class Sdf06Executor {
|
||||||
Instant updatedTime = Instant.now();
|
Instant updatedTime = Instant.now();
|
||||||
if (gatewayMsg.isApproved()) {
|
if (gatewayMsg.isApproved()) {
|
||||||
processedApproved(stmt, sdf06, updatedTime);
|
processedApproved(stmt, sdf06, updatedTime);
|
||||||
dmiService.setProc(searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null),
|
dmiService.setProcContract(searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null),
|
||||||
CurrencyCode.RUB.getKey(),
|
CurrencyCode.RUB.getKey(),
|
||||||
safeBD(sdf06.getSum()),
|
safeBD(sdf06.getSum()),
|
||||||
sdf06.getNumber().toString()
|
sdf06.getNumber().toString()
|
||||||
|
|
|
||||||
|
|
@ -50,7 +50,7 @@ public class Sdf55Executor {
|
||||||
this.assets = assets;
|
this.assets = assets;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void execute(BaseRequest<StatementRequest> systemRequest) {
|
public synchronized void execute(BaseRequest<StatementRequest> systemRequest) {
|
||||||
StatementRequest payload = systemRequest.getRequestPayload();
|
StatementRequest payload = systemRequest.getRequestPayload();
|
||||||
log.debug("executing {}.generationId={}", payload.getTable(), payload.getGroupId());
|
log.debug("executing {}.generationId={}", payload.getTable(), payload.getGroupId());
|
||||||
IValidator validator = validation.apply(payload);
|
IValidator validator = validation.apply(payload);
|
||||||
|
|
@ -68,12 +68,10 @@ public class Sdf55Executor {
|
||||||
pmt.setTransactionStatus(TransactionStatus.ok.getKey());
|
pmt.setTransactionStatus(TransactionStatus.ok.getKey());
|
||||||
pmt.setUpdated(Instant.now());
|
pmt.setUpdated(Instant.now());
|
||||||
pmtImdg.update(pmt);
|
pmtImdg.update(pmt);
|
||||||
if (tcr != null) {
|
dmiService.setContract(tcr.getId(), CurrencyCode.RUB.getKey(), sDf54.getDocnm_ref(), sDf55.getDocnmprev());
|
||||||
dmiService.setOk(tcr.getId(), CurrencyCode.RUB.getKey(), sDf54.getDocnm_ref());
|
assets.process(assetTrio.a__b(),
|
||||||
assets.process(assetTrio.a__b(),
|
assetTrio.a__t(),
|
||||||
assetTrio.a__t(),
|
assetTrio.a__f(),
|
||||||
assetTrio.a__f(),
|
BigDecimal.ZERO);
|
||||||
BigDecimal.ZERO);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -25,6 +25,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationLi
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.clearing.service.AnltSearcher;
|
import ru.spcex.clearing.service.AnltSearcher;
|
||||||
|
import ru.spcex.clearing.service.AssetTrio;
|
||||||
import ru.spcex.clearing.service.LoggingService;
|
import ru.spcex.clearing.service.LoggingService;
|
||||||
import ru.spcex.clearing.service.builder.RegistryBuilder;
|
import ru.spcex.clearing.service.builder.RegistryBuilder;
|
||||||
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
|
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
|
||||||
|
|
@ -56,7 +57,6 @@ import java.util.Collection;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
import java.util.concurrent.atomic.AtomicReference;
|
|
||||||
import java.util.function.Consumer;
|
import java.util.function.Consumer;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
import java.util.regex.Pattern;
|
import java.util.regex.Pattern;
|
||||||
|
|
@ -197,7 +197,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
statementDeb.map(SpcexObjectBase::getId).orElse(null),
|
statementDeb.map(SpcexObjectBase::getId).orElse(null),
|
||||||
statementCred.map(SpcexObjectBase::getId).orElse(null));
|
statementCred.map(SpcexObjectBase::getId).orElse(null));
|
||||||
|
|
||||||
Function<StmtCmpAcc, Optional<Long>> createRegistryIfNeeded = stmtCmpAcc -> {
|
Function<StmtCmpAcc, Optional<AssetTrio>> createRegistryIfNeeded = stmtCmpAcc -> {
|
||||||
Statement stmt = stmtCmpAcc.statement();
|
Statement stmt = stmtCmpAcc.statement();
|
||||||
Company company = stmtCmpAcc.company();
|
Company company = stmtCmpAcc.company();
|
||||||
Account account = stmtCmpAcc.account();
|
Account account = stmtCmpAcc.account();
|
||||||
|
|
@ -208,8 +208,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
.companyId(company.getId()) //раньше было stmt.getAddresseeId()
|
.companyId(company.getId()) //раньше было stmt.getAddresseeId()
|
||||||
.accountId(account.getId()) //раньше было stmt.getAccountId()
|
.accountId(account.getId()) //раньше было stmt.getAccountId()
|
||||||
.find();
|
.find();
|
||||||
AtomicReference<Registry> amtCached = new AtomicReference<>();
|
AssetTrio asts = amtFound.map(rgs -> {
|
||||||
amtFound.ifPresentOrElse(rgs -> {
|
|
||||||
log.debug("stmt.id={}, found AM*T.id={}, updating...", stmt.getId(), rgs.getId());
|
log.debug("stmt.id={}, found AM*T.id={}, updating...", stmt.getId(), rgs.getId());
|
||||||
updateReg(stmt, rgs);
|
updateReg(stmt, rgs);
|
||||||
//добавил создание если не найдены
|
//добавил создание если не найдены
|
||||||
|
|
@ -239,8 +238,8 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
rgs.getBalance(),
|
rgs.getBalance(),
|
||||||
rgs.getDebit(),
|
rgs.getDebit(),
|
||||||
rgs.getCredit());
|
rgs.getCredit());
|
||||||
amtCached.set(rgs);
|
return new AssetTrio(registryUnitF, registryUnitB, rgs);
|
||||||
}, () -> {
|
}).orElseGet(() -> {
|
||||||
Registry registry = RegistryBuilder.builder(imdgProvider)
|
Registry registry = RegistryBuilder.builder(imdgProvider)
|
||||||
.statement(stmt)
|
.statement(stmt)
|
||||||
.company(company)
|
.company(company)
|
||||||
|
|
@ -259,7 +258,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
registry.getId(),
|
registry.getId(),
|
||||||
registryF.getId(),
|
registryF.getId(),
|
||||||
registryB.getId());
|
registryB.getId());
|
||||||
amtCached.set(registry);
|
return new AssetTrio(registryF, registryB, registry);
|
||||||
});
|
});
|
||||||
InOutDirection direction = IEnumKey.getEnumByKey(InOutDirection.class, stmt.getInOutDirection());
|
InOutDirection direction = IEnumKey.getEnumByKey(InOutDirection.class, stmt.getInOutDirection());
|
||||||
boolean dm_tWasCreated = false;
|
boolean dm_tWasCreated = false;
|
||||||
|
|
@ -334,10 +333,9 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
kafka.sendRequestToQueue(Consts.ASSET_OPERATION,
|
kafka.sendRequestToQueue(Consts.ASSET_OPERATION,
|
||||||
gatewayRequest(stmt,
|
gatewayRequest(stmt,
|
||||||
company.getTradingCode(),
|
company.getTradingCode(),
|
||||||
amtCached.get().getTradingClearingRegistry()));
|
asts.a__t().getTradingClearingRegistry()));
|
||||||
|
|
||||||
}
|
}
|
||||||
return Optional.ofNullable(amtCached.get().getTradingClearingRegistryId());
|
return Optional.of(asts);
|
||||||
} else {
|
} else {
|
||||||
log.debug("statement.id={} activness validation failed {}", stmt.getId(), messageResolver.resolve(err.get()));
|
log.debug("statement.id={} activness validation failed {}", stmt.getId(), messageResolver.resolve(err.get()));
|
||||||
stmt.setErrorCodeId(err.get().getSubject().getId()); // fixme ErrorText insert
|
stmt.setErrorCodeId(err.get().getSubject().getId()); // fixme ErrorText insert
|
||||||
|
|
@ -350,7 +348,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
Statement stmt = stmtCmpAcc.statement();
|
Statement stmt = stmtCmpAcc.statement();
|
||||||
Company cmp = stmtCmpAcc.company();
|
Company cmp = stmtCmpAcc.company();
|
||||||
Account acc = stmtCmpAcc.account();
|
Account acc = stmtCmpAcc.account();
|
||||||
Optional<Long> tcrId = createRegistryIfNeeded.apply(stmtCmpAcc);
|
Optional<AssetTrio> asts = createRegistryIfNeeded.apply(stmtCmpAcc);
|
||||||
if (accountIsAnlt(acc)) {
|
if (accountIsAnlt(acc)) {
|
||||||
AnltSearcher.AnltSearch searchResult = anltSearcher.loadByAnlt(stmt.getComment());
|
AnltSearcher.AnltSearch searchResult = anltSearcher.loadByAnlt(stmt.getComment());
|
||||||
if (searchResult.isFound()) {
|
if (searchResult.isFound()) {
|
||||||
|
|
@ -359,11 +357,12 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
stmt.getComment(),
|
stmt.getComment(),
|
||||||
searchResult.getAccount().getId(),
|
searchResult.getAccount().getId(),
|
||||||
searchResult.getCompany().getId());
|
searchResult.getCompany().getId());
|
||||||
createRegistryIfNeeded.apply(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount()));
|
asts = createRegistryIfNeeded.apply(new StmtCmpAcc(stmt, searchResult.getCompany(), searchResult.getAccount()));
|
||||||
if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
|
if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
|
||||||
dmiService.setOk(searchResult.getTcr().getId(),
|
dmiService.setOk(searchResult.getTcr().getId(),
|
||||||
CurrencyCode.RUB.getKey(),
|
CurrencyCode.RUB.getKey(),
|
||||||
sdf57.getDbfId().toString());
|
sdf57.getDbfId().toString());
|
||||||
|
asts.ifPresent(trio -> assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO));
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
log.debug("stmt.id={} comment='{}' error: {}. Operating through DMAU registry",
|
log.debug("stmt.id={} comment='{}' error: {}. Operating through DMAU registry",
|
||||||
|
|
@ -396,10 +395,11 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
//DMAU подразумевает увеличение остатка поля balance и credit registry.code=DMAU
|
//DMAU подразумевает увеличение остатка поля balance и credit registry.code=DMAU
|
||||||
}
|
}
|
||||||
} else if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
|
} else if (OperationStatus.Executed.equalsByKey(stmt.getOperationStatus())) {
|
||||||
if (tcrId.isPresent()) {
|
if (asts.isPresent()) {
|
||||||
dmiService.setOk(tcrId.get(),
|
dmiService.setOk(asts.get().a__t().getTradingClearingRegistryId(),
|
||||||
CurrencyCode.RUB.getKey(),
|
CurrencyCode.RUB.getKey(),
|
||||||
sdf57.getDbfId().toString());
|
sdf57.getDbfId().toString());
|
||||||
|
asts.ifPresent(trio -> assets.process(trio.a__b(), trio.a__t(), trio.a__f(), BigDecimal.ZERO));
|
||||||
} else {
|
} else {
|
||||||
log.error("stmt.id={} operationStatus={} but couldn't extract TCR.id for DM*I update",
|
log.error("stmt.id={} operationStatus={} but couldn't extract TCR.id for DM*I update",
|
||||||
stmt.getId(),
|
stmt.getId(),
|
||||||
|
|
@ -679,7 +679,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
|
|
||||||
private AssetOperationListRequest gatewayRequest(Statement stmt, String tradingCode, String tcrCode) {
|
private AssetOperationListRequest gatewayRequest(Statement stmt, String tradingCode, String tcrCode) {
|
||||||
AssetOperationListRequest gatewayRequest = new AssetOperationListRequest();
|
AssetOperationListRequest gatewayRequest = new AssetOperationListRequest();
|
||||||
AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart(stmt, tradingCode, tcrCode);
|
AssetOperationRequest req = GatewayRequestCreator.from(stmt, tradingCode, tcrCode);
|
||||||
Session existActiveSession = sessionImdg.getFirstObjectByFieldValues(Map.of(
|
Session existActiveSession = sessionImdg.getFirstObjectByFieldValues(Map.of(
|
||||||
"workflowStatus", SessionStatus.ACTV.getKey()
|
"workflowStatus", SessionStatus.ACTV.getKey()
|
||||||
));
|
));
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,7 @@
|
||||||
package ru.spcex.clearing.service.integration;
|
package ru.spcex.clearing.service.integration;
|
||||||
|
|
||||||
|
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||||
|
import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
import ru.clearing.classes.statics.data.statement.Statement;
|
import ru.clearing.classes.statics.data.statement.Statement;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
|
||||||
import ru.spcex.platform.enumeration.CurrencyCode;
|
import ru.spcex.platform.enumeration.CurrencyCode;
|
||||||
|
|
@ -8,7 +10,7 @@ import ru.spcex.platform.enumeration.InOutDirection;
|
||||||
import java.math.BigDecimal;
|
import java.math.BigDecimal;
|
||||||
|
|
||||||
public class GatewayRequestCreator {
|
public class GatewayRequestCreator {
|
||||||
public static AssetOperationRequest gatewayRequestPart(Statement stmt, String tradingCode, String tcrCode) {
|
public static AssetOperationRequest from(Statement stmt, String tradingCode, String tcrCode) {
|
||||||
AssetOperationRequest req = new AssetOperationRequest();
|
AssetOperationRequest req = new AssetOperationRequest();
|
||||||
req.setEntityId(stmt.getId());
|
req.setEntityId(stmt.getId());
|
||||||
req.setAmount(stmt.getAmount());
|
req.setAmount(stmt.getAmount());
|
||||||
|
|
@ -19,12 +21,30 @@ public class GatewayRequestCreator {
|
||||||
return req;
|
return req;
|
||||||
}
|
}
|
||||||
|
|
||||||
public static AssetOperationRequest gatewayRequestPart(Long entityId,
|
public static AssetOperationRequest from(Registry om_t) {
|
||||||
Long sessionId,
|
return from(om_t.getId(),
|
||||||
BigDecimal amount,
|
null,
|
||||||
InOutDirection direction,
|
om_t.getBalance(),
|
||||||
String tradingCode,
|
InOutDirection.out,
|
||||||
String tcrCode) {
|
om_t.getTradingCode(),
|
||||||
|
om_t.getTradingClearingRegistry());
|
||||||
|
}
|
||||||
|
|
||||||
|
public static AssetOperationRequest from(PaymentInstruction pmt, String tradingCode, String tcr) {
|
||||||
|
return from(pmt.getId(),
|
||||||
|
null,
|
||||||
|
pmt.getCreditLeg_amount(),
|
||||||
|
InOutDirection.out,
|
||||||
|
tradingCode,
|
||||||
|
tcr);
|
||||||
|
}
|
||||||
|
|
||||||
|
public static AssetOperationRequest from(Long entityId,
|
||||||
|
Long sessionId,
|
||||||
|
BigDecimal amount,
|
||||||
|
InOutDirection direction,
|
||||||
|
String tradingCode,
|
||||||
|
String tcrCode) {
|
||||||
AssetOperationRequest req = new AssetOperationRequest();
|
AssetOperationRequest req = new AssetOperationRequest();
|
||||||
req.setSessionId(sessionId);
|
req.setSessionId(sessionId);
|
||||||
req.setEntityId(entityId);
|
req.setEntityId(entityId);
|
||||||
|
|
|
||||||
|
|
@ -14,8 +14,10 @@ import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf54;
|
import ru.clearing.classes.statics.data.sdf.SDf54;
|
||||||
import ru.spcex.clearing.error.ClearingError;
|
import ru.spcex.clearing.error.ClearingError;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.notification.NotificationSender;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
|
|
@ -24,14 +26,15 @@ import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.clearing.service.AssetTrio;
|
import ru.spcex.clearing.service.AssetTrio;
|
||||||
import ru.spcex.clearing.service.Sdf54Creator;
|
import ru.spcex.clearing.service.Sdf54Creator;
|
||||||
import ru.spcex.clearing.service.builder.PaymentInstructionBuilderV2;
|
import ru.spcex.clearing.service.builder.PaymentInstructionBuilderV2;
|
||||||
|
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
|
||||||
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
|
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
|
||||||
import ru.spcex.clearing.service.registry.DmiService;
|
import ru.spcex.clearing.service.registry.DmiService;
|
||||||
import ru.spcex.clearing.service.registry.RegistryManager;
|
import ru.spcex.clearing.service.registry.RegistryManager;
|
||||||
|
import ru.spcex.clearing.service.schedule.TradingTimeService;
|
||||||
import ru.spcex.clearing.service.validation.ValidationStored;
|
import ru.spcex.clearing.service.validation.ValidationStored;
|
||||||
|
import ru.spcex.clearing.session.stage.impl.GatewayRequester;
|
||||||
import ru.spcex.clearing.util.security.UserRoleVerification;
|
import ru.spcex.clearing.util.security.UserRoleVerification;
|
||||||
import ru.spcex.platform.enumeration.CurrencyCode;
|
import ru.spcex.platform.enumeration.*;
|
||||||
import ru.spcex.platform.enumeration.TransactionStatus;
|
|
||||||
import ru.spcex.platform.enumeration.UserRole;
|
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgId;
|
import ru.spcex.platform.imdg.api.ImdgId;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
@ -44,6 +47,7 @@ import java.math.BigDecimal;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
|
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
|
||||||
|
|
||||||
|
|
@ -67,12 +71,15 @@ public class PaymentInstructionOutboundService {
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
private final AssetTBFProcessing assetsMng;
|
private final AssetTBFProcessing assetsMng;
|
||||||
private final DmiService dmiService;
|
private final DmiService dmiService;
|
||||||
|
private final TradingTimeService time;
|
||||||
|
private final GatewayRequester gateway;
|
||||||
|
private final NotificationSender notification;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public PaymentInstructionOutboundService(UserRoleVerification rights,
|
public PaymentInstructionOutboundService(UserRoleVerification rights,
|
||||||
IMessageResolver msgs,
|
IMessageResolver msgs,
|
||||||
ImdgProvider imdgProvider,
|
ImdgProvider imdgProvider,
|
||||||
Function<PIClearingOutbondActionNewRequest, IValidator> validation, RegistryManager rgsMng, KafkaSender kafkaSender, AssetTBFProcessing assetsMng, DmiService dmiService) {
|
Function<PIClearingOutbondActionNewRequest, IValidator> validation, RegistryManager rgsMng, KafkaSender kafkaSender, AssetTBFProcessing assetsMng, DmiService dmiService, TradingTimeService time, GatewayRequester gateway, NotificationSender notification) {
|
||||||
this.imdgProvider = imdgProvider;
|
this.imdgProvider = imdgProvider;
|
||||||
this.rights = rights;
|
this.rights = rights;
|
||||||
this.msgs = msgs;
|
this.msgs = msgs;
|
||||||
|
|
@ -90,6 +97,10 @@ public class PaymentInstructionOutboundService {
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
this.assetsMng = assetsMng;
|
this.assetsMng = assetsMng;
|
||||||
this.dmiService = dmiService;
|
this.dmiService = dmiService;
|
||||||
|
this.time = time;
|
||||||
|
this.gateway = gateway;
|
||||||
|
this.gateway.setName("PaymentInstructionOutboundService|PI");
|
||||||
|
this.notification = notification;
|
||||||
}
|
}
|
||||||
|
|
||||||
public RequestInfoUpdate sendOutBoundPayment(BaseRequest<PIClearingOutbondActionNewRequest> req) {
|
public RequestInfoUpdate sendOutBoundPayment(BaseRequest<PIClearingOutbondActionNewRequest> req) {
|
||||||
|
|
@ -106,14 +117,13 @@ public class PaymentInstructionOutboundService {
|
||||||
}
|
}
|
||||||
Account accCred = validator.getStored(ValidationStored.PIOutboundAccountCred);
|
Account accCred = validator.getStored(ValidationStored.PIOutboundAccountCred);
|
||||||
Account accDeb = validator.getStored(ValidationStored.PIOutboundAccountDeb);
|
Account accDeb = validator.getStored(ValidationStored.PIOutboundAccountDeb);
|
||||||
Company addressee = cmpImdg.getSingleObjectByID(payload.getAddresseeId());
|
Company addressee = validator.getStored(ValidationStored.PIOOutboundAddressee);
|
||||||
|
TradingClearingRegistry tcr = validator.getStored(ValidationStored.PIOutboundTcr);
|
||||||
Company sender = cmpImdg.getSingleObjectByID(payload.getSenderId());
|
Company sender = cmpImdg.getSingleObjectByID(payload.getSenderId());
|
||||||
BigDecimal amount = safeBD(payload.getCreditLeg_amount());
|
BigDecimal amount = safeBD(payload.getCreditLeg_amount());
|
||||||
log.debug("all checks passed, accCred.id={}, accDeb.id={}, addressee.id={}, sender.id={}, amount: {}",
|
log.debug("all checks passed, accCred.id={}, accDeb.id={}, addressee.id={}, sender.id={}, amount: {}",
|
||||||
accCred.getId(), accDeb.getId(), addressee.getId(), sender.getId(), amount);
|
accCred.getId(), accDeb.getId(), addressee.getId(), sender.getId(), amount);
|
||||||
|
|
||||||
TradingClearingRegistry tcr = tcrImdg.getFirstObjectBySQL("moneyAccountId = '%d' and companyId = %d"
|
|
||||||
.formatted(accCred.getId(), addressee.getId()));
|
|
||||||
String purpose = tcr != null ? "Вывод средств по ТКР " + tcr.getCode() + "." : "Вывод средств.";
|
String purpose = tcr != null ? "Вывод средств по ТКР " + tcr.getCode() + "." : "Вывод средств.";
|
||||||
if (payload.getPaymentPurpose() != null) {
|
if (payload.getPaymentPurpose() != null) {
|
||||||
purpose += " " + payload.getPaymentPurpose();
|
purpose += " " + payload.getPaymentPurpose();
|
||||||
|
|
@ -141,7 +151,7 @@ public class PaymentInstructionOutboundService {
|
||||||
pmt.setCreditLeg_securityId(currency.getId());
|
pmt.setCreditLeg_securityId(currency.getId());
|
||||||
pmt.setDebitLeg_securityId(currency.getId());
|
pmt.setDebitLeg_securityId(currency.getId());
|
||||||
}
|
}
|
||||||
pmt.setTransactionStatus(TransactionStatus.exec.getKey());
|
pmt.setTransactionStatus(TransactionStatus.stld.getKey());
|
||||||
//fixme pmt.setDocumentNumber(registry.getSecurityId());
|
//fixme pmt.setDocumentNumber(registry.getSecurityId());
|
||||||
|
|
||||||
pmtImdg.insert(pmt);
|
pmtImdg.insert(pmt);
|
||||||
|
|
@ -162,12 +172,28 @@ public class PaymentInstructionOutboundService {
|
||||||
sDf54.setGenerationId(pmt.getId());
|
sDf54.setGenerationId(pmt.getId());
|
||||||
sdf54Imdg.insert(sDf54);
|
sdf54Imdg.insert(sDf54);
|
||||||
log.debug("new sdf54.id: {}", sDf54.getId());
|
log.debug("new sdf54.id: {}", sDf54.getId());
|
||||||
if (tcr != null) {
|
dmiService.setProcComment(tcr.getId(),
|
||||||
dmiService.setProc(tcr.getId(),
|
CurrencyCode.RUB.getKey(),
|
||||||
CurrencyCode.RUB.getKey(),
|
safeBD(amount).negate(),
|
||||||
safeBD(amount).negate(),
|
sDf54.getDocnm_ref());
|
||||||
sDf54.getDocnm_ref());
|
assetsMng.process(assets.get().a__b(), assets.get().a__t(), assets.get().a__f(), BigDecimal.ZERO);
|
||||||
assetsMng.process(assets.get().a__b(), assets.get().a__t(), assets.get().a__f(), BigDecimal.ZERO);
|
if (time.isTradingTime()) {
|
||||||
|
Supplier<AssetOperationRequest> builder = () -> GatewayRequestCreator.from(pmt,
|
||||||
|
assets.get().a__b().getTradingCode(),
|
||||||
|
tcr.getCode());
|
||||||
|
Optional<Boolean> gatewayOk = gateway.gatewayRequestAndWait(builder);
|
||||||
|
if (gatewayOk.isEmpty() || gatewayOk.get().equals(Boolean.FALSE)) {
|
||||||
|
log.error("sdf54.id={} pmt.id={} gateway {}", sDf54.getId(),
|
||||||
|
pmt.getId(),
|
||||||
|
gatewayOk.isPresent() ? "failed" : "timeout");
|
||||||
|
notification.sendNotification(ObjectType.gateway,
|
||||||
|
gatewayOk.isPresent() ? "gateway not approved" : "gateway timeout",
|
||||||
|
Priority.HIGH);
|
||||||
|
pmt.setTransactionStatus(TransactionStatus.fail.getKey());
|
||||||
|
pmt.setUpdated(now);
|
||||||
|
pmtImdg.update(pmt);
|
||||||
|
return error(req.getId(), ClearingError.GeneralError);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
SdfClearingRequest exp = new SdfClearingRequest();
|
SdfClearingRequest exp = new SdfClearingRequest();
|
||||||
exp.setGroupId(sDf54.getGenerationId());
|
exp.setGroupId(sDf54.getGenerationId());
|
||||||
|
|
|
||||||
|
|
@ -31,23 +31,76 @@ public class DmiService {
|
||||||
this.rgsMng = rgsMng;
|
this.rgsMng = rgsMng;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setProc(Long tcrId,
|
/**
|
||||||
String securitySymbol,
|
* ищем D**I по tcrId/security/contract
|
||||||
BigDecimal balance,
|
* обновляем либо создаем
|
||||||
String contract) {
|
*/
|
||||||
findDmi(tcrId, securitySymbol, contract.toString()).ifPresentOrElse(d__i -> {
|
public void setProcContract(Long tcrId,
|
||||||
d__i.setRegistryStatus(RegistryStatus.PROC.getKey());
|
String securitySymbol,
|
||||||
|
BigDecimal balance,
|
||||||
|
String contract) {
|
||||||
|
findDmiByContract(tcrId, securitySymbol, contract).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(tcrId, securitySymbol, balance)
|
||||||
|
.ifPresent(d__i -> {
|
||||||
|
d__i.setContract(contract);
|
||||||
|
rgsImdg.insert(d__i);
|
||||||
|
log.debug("created {}.id={}", d__i.getRegistryCode(), d__i.getId());
|
||||||
|
})
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* ищем D**I по tcrId/security/comment
|
||||||
|
* обновляем либо создаем
|
||||||
|
*/
|
||||||
|
public void setProcComment(Long tcrId,
|
||||||
|
String securitySymbol,
|
||||||
|
BigDecimal balance,
|
||||||
|
String comment) {
|
||||||
|
findDmiByComment(tcrId, securitySymbol, comment).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(tcrId, securitySymbol, balance)
|
||||||
|
.ifPresent(d__i -> {
|
||||||
|
d__i.setComment(comment);
|
||||||
|
rgsImdg.insert(d__i);
|
||||||
|
log.debug("created {}.id={}", d__i.getRegistryCode(), d__i.getId());
|
||||||
|
})
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* ищем D**I по tcrId/security/comment
|
||||||
|
* добавляем поле contract
|
||||||
|
*/
|
||||||
|
public void setContract(Long tcrId,
|
||||||
|
String securitySymbol,
|
||||||
|
String comment,
|
||||||
|
String contract) {
|
||||||
|
findDmiByComment(tcrId, securitySymbol, comment).ifPresentOrElse(d__i -> {
|
||||||
|
d__i.setContract(contract);
|
||||||
d__i.setUpdated(Instant.now());
|
d__i.setUpdated(Instant.now());
|
||||||
log.trace("updating D**I status to {}", d__i.getRegistryStatus());
|
|
||||||
rgsImdg.update(d__i);
|
rgsImdg.update(d__i);
|
||||||
}, () -> createDmi(tcrId, securitySymbol, balance, contract.toString()));
|
log.trace("set {}.id={} contract {}", d__i.getRegistryCode(), d__i.getId(), contract);
|
||||||
|
}, () -> log.trace("didn't find DM*I by tcr.id={} securitySymbol={} comment={}",
|
||||||
|
tcrId, securitySymbol, comment));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
/**
|
||||||
|
* ищем D**I по tcrId/security/contract
|
||||||
|
* задаем статус ОК
|
||||||
|
*/
|
||||||
public void setOk(Long tcrId,
|
public void setOk(Long tcrId,
|
||||||
String securitySymbol,
|
String securitySymbol,
|
||||||
String contract) {
|
String contract) {
|
||||||
findDmi(tcrId, securitySymbol, contract).ifPresentOrElse(d__i -> {
|
findDmiByContract(tcrId, securitySymbol, contract).ifPresentOrElse(d__i -> {
|
||||||
d__i.setRegistryStatus(RegistryStatus.OK.getKey());
|
d__i.setRegistryStatus(RegistryStatus.OK.getKey());
|
||||||
d__i.setUpdated(Instant.now());
|
d__i.setUpdated(Instant.now());
|
||||||
log.trace("updating {} status to {}", d__i.getRegistryCode(), d__i.getRegistryStatus());
|
log.trace("updating {} status to {}", d__i.getRegistryCode(), d__i.getRegistryStatus());
|
||||||
|
|
@ -58,7 +111,7 @@ public class DmiService {
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private Optional<Registry> findDmi(Long tcrId, String securitySymbol, String contract) {
|
private Optional<Registry> findDmiByContract(Long tcrId, String securitySymbol, String contract) {
|
||||||
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
|
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
|
||||||
RegistryTradingParams rgsCde = RegistryTradingParams.D__I;
|
RegistryTradingParams rgsCde = RegistryTradingParams.D__I;
|
||||||
ImdgPredicate prdct = pb.and(
|
ImdgPredicate prdct = pb.and(
|
||||||
|
|
@ -72,10 +125,23 @@ public class DmiService {
|
||||||
return Optional.ofNullable(d__i);
|
return Optional.ofNullable(d__i);
|
||||||
}
|
}
|
||||||
|
|
||||||
private void createDmi(Long tcrId,
|
private Optional<Registry> findDmiByComment(Long tcrId, String securitySymbol, String comment) {
|
||||||
|
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
|
||||||
|
RegistryTradingParams rgsCde = RegistryTradingParams.D__I;
|
||||||
|
ImdgPredicate prdct = pb.and(
|
||||||
|
pb.sql(RegistryCodeSqlBuilder.getInstance(rgsCde).build()),
|
||||||
|
pb.equals("tradingClearingRegistryId", tcrId),
|
||||||
|
pb.equals("securitySymbol", securitySymbol),
|
||||||
|
pb.equals("comment", comment)
|
||||||
|
);
|
||||||
|
log.trace("searching D**I by {}", prdct);
|
||||||
|
Registry d__i = rgsImdg.getFirstObjectByPredicate(prdct);
|
||||||
|
return Optional.ofNullable(d__i);
|
||||||
|
}
|
||||||
|
|
||||||
|
private Optional<Registry> createDmi(Long tcrId,
|
||||||
String securitySymbol,
|
String securitySymbol,
|
||||||
BigDecimal summ,
|
BigDecimal summ) {
|
||||||
String outDocument) {
|
|
||||||
RegistryTradingParams rgsCde = RegistryTradingParams.AM_T;
|
RegistryTradingParams rgsCde = RegistryTradingParams.AM_T;
|
||||||
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
|
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
|
||||||
ImdgPredicate prdct = pb.and(
|
ImdgPredicate prdct = pb.and(
|
||||||
|
|
@ -87,7 +153,7 @@ public class DmiService {
|
||||||
if (am_t == null) {
|
if (am_t == null) {
|
||||||
log.error("AM*T not found for tcr.id={} security.symbol={}",
|
log.error("AM*T not found for tcr.id={} security.symbol={}",
|
||||||
tcrId, securitySymbol);
|
tcrId, securitySymbol);
|
||||||
return;
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
Registry rgsD = am_t.clone();
|
Registry rgsD = am_t.clone();
|
||||||
rgsD.setRegistryDesignation(RegistryDesignation.D.getKey());
|
rgsD.setRegistryDesignation(RegistryDesignation.D.getKey());
|
||||||
|
|
@ -100,8 +166,6 @@ public class DmiService {
|
||||||
rgsD.setDiffBalance(BigDecimal.ZERO);
|
rgsD.setDiffBalance(BigDecimal.ZERO);
|
||||||
rgsD.setCheckBalance(BigDecimal.ZERO);
|
rgsD.setCheckBalance(BigDecimal.ZERO);
|
||||||
rgsD.setRegistryStatus(RegistryStatus.PROC.getKey());
|
rgsD.setRegistryStatus(RegistryStatus.PROC.getKey());
|
||||||
rgsD.setContract(outDocument);
|
return Optional.of(rgsD);
|
||||||
rgsImdg.insert(rgsD);
|
|
||||||
log.debug("created {}.id={}", rgsD.getRegistryCode(), rgsD.getId());
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,8 @@ package ru.spcex.clearing.service.validation;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import ru.clearing.classes.statics.data.account.Account;
|
import ru.clearing.classes.statics.data.account.Account;
|
||||||
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
|
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
import ru.spcex.clearing.error.ClearingError;
|
import ru.spcex.clearing.error.ClearingError;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
||||||
|
|
@ -80,6 +82,34 @@ public enum PaymentOutboundValidationRule implements IValidationRule<ImdgValidat
|
||||||
return empty();
|
return empty();
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
AddresseePresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<PIClearingOutbondActionNewRequest> context) {
|
||||||
|
PIClearingOutbondActionNewRequest validatedObject = context.getValidatedObject();
|
||||||
|
Imdg<Company> cmpImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
|
Company addressee = cmpImdg.getSingleObjectByID(validatedObject.getAddresseeId());
|
||||||
|
if (addressee == null) {
|
||||||
|
return of(ClearingError.WrongFieldValue, "addresseeId");
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.PIOOutboundAddressee, addressee);
|
||||||
|
return Optional.empty();
|
||||||
|
}
|
||||||
|
},
|
||||||
|
TcrPresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<PIClearingOutbondActionNewRequest> context) {
|
||||||
|
Account accCred = context.getStoredObject(ValidationStored.PIOutboundAccountCred);
|
||||||
|
Company addressee = context.getStoredObject(ValidationStored.PIOOutboundAddressee);
|
||||||
|
Imdg<TradingClearingRegistry> tcrImdg = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
|
||||||
|
TradingClearingRegistry tcr = tcrImdg.getFirstObjectBySQL("moneyAccountId = '%d' and companyId = %d"
|
||||||
|
.formatted(accCred.getId(), addressee.getId()));
|
||||||
|
if (tcr == null) {
|
||||||
|
return of(ClearingError.RecordNotFound, "TCR");
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.PIOutboundTcr, tcr);
|
||||||
|
return Optional.empty();
|
||||||
|
}
|
||||||
|
}
|
||||||
;
|
;
|
||||||
private final static Logger log = LoggerFactory.getLogger(PaymentOutboundValidationRule.class);
|
private final static Logger log = LoggerFactory.getLogger(PaymentOutboundValidationRule.class);
|
||||||
@Override
|
@Override
|
||||||
|
|
|
||||||
|
|
@ -77,13 +77,18 @@ public enum Sdf55ValidationRule implements IValidationRule<ImdgValidationContext
|
||||||
PaymentInstruction pmt = context.getStoredObject(ValidationStored.Sdf55PaymentInstruction);
|
PaymentInstruction pmt = context.getStoredObject(ValidationStored.Sdf55PaymentInstruction);
|
||||||
Long accCredId = pmt.getCreditLeg_accountId();
|
Long accCredId = pmt.getCreditLeg_accountId();
|
||||||
Long addresseeId = pmt.getSenderId();
|
Long addresseeId = pmt.getSenderId();
|
||||||
if (accCredId == null || addresseeId == null) return empty();
|
if (accCredId == null || addresseeId == null) {
|
||||||
|
return of(ClearingError.RecordNotFound, "TCR - no PI(id=%d)#CreditLeg_accountId/SenderId"
|
||||||
|
.formatted(pmt.getId()));
|
||||||
|
}
|
||||||
TradingClearingRegistry tcr = tcrImdg.getFirstObjectBySQL(
|
TradingClearingRegistry tcr = tcrImdg.getFirstObjectBySQL(
|
||||||
"moneyAccountId = '%d' and companyId = %d".formatted(accCredId, addresseeId)
|
"moneyAccountId = '%d' and companyId = %d".formatted(accCredId, addresseeId)
|
||||||
);
|
);
|
||||||
if (tcr != null) {
|
if (tcr == null) {
|
||||||
context.storeObject(ValidationStored.Sdf55Tcr, tcr);
|
return of(ClearingError.RecordNotFound, "TCR moneyAccountId = '%d' and companyId = %d"
|
||||||
|
.formatted(accCredId, addresseeId));
|
||||||
}
|
}
|
||||||
|
context.storeObject(ValidationStored.Sdf55Tcr, tcr);
|
||||||
return empty();
|
return empty();
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|
|
||||||
|
|
@ -14,7 +14,7 @@ public enum ValidationStored {
|
||||||
|
|
||||||
Sdf21Account,
|
Sdf21Account,
|
||||||
|
|
||||||
PIOutboundAccountDeb, PIOutboundAccountCred,
|
PIOutboundAccountDeb, PIOutboundAccountCred, PIOOutboundAddressee, PIOutboundTcr,
|
||||||
|
|
||||||
Sdf06Company, Sdf06Account, Sdf06Tcr,
|
Sdf06Company, Sdf06Account, Sdf06Tcr,
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -9,7 +9,6 @@ import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
||||||
import org.springframework.context.annotation.Scope;
|
import org.springframework.context.annotation.Scope;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
import ru.clearing.classes.statics.data.registry.Registry;
|
|
||||||
import ru.spcex.clearing.config.element.ClearingServiceSettings;
|
import ru.spcex.clearing.config.element.ClearingServiceSettings;
|
||||||
import ru.spcex.clearing.config.element.SessionStageSettings;
|
import ru.spcex.clearing.config.element.SessionStageSettings;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
|
@ -20,8 +19,6 @@ import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApp
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
|
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
|
|
||||||
import ru.spcex.platform.enumeration.InOutDirection;
|
|
||||||
import ru.spcex.platform.utils.log.ExceptionUtils;
|
import ru.spcex.platform.utils.log.ExceptionUtils;
|
||||||
|
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
|
|
@ -40,6 +37,7 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
private final Integer GATEWAY_TIMEOUT;
|
private final Integer GATEWAY_TIMEOUT;
|
||||||
|
private String name;
|
||||||
|
|
||||||
private final Lock gatewayLock = new ReentrantLock();
|
private final Lock gatewayLock = new ReentrantLock();
|
||||||
private final Condition gatewayCondition = gatewayLock.newCondition();
|
private final Condition gatewayCondition = gatewayLock.newCondition();
|
||||||
|
|
@ -65,30 +63,26 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
public Optional<Boolean> gatewayRequestAndWait(Registry om_t) {
|
public Optional<Boolean> gatewayRequestAndWait(Supplier<AssetOperationRequest> reqBuilder) {
|
||||||
SingleAssetResponse responseReceived = null;
|
SingleAssetResponse responseReceived = null;
|
||||||
|
Long entityId = null;
|
||||||
gatewayLock.lock();
|
gatewayLock.lock();
|
||||||
try {
|
try {
|
||||||
if (this.omtId != null) {
|
if (this.omtId != null) {
|
||||||
log.error("om*t.id={} already sent to gateway", this.omtId);
|
log.error("{}.id={} already sent to gateway", getName(), this.omtId);
|
||||||
throw new IllegalStateException(" forbid reusing GatewayRequester instance in multiple threads");
|
throw new IllegalStateException(" forbid reusing GatewayRequester instance in multiple threads");
|
||||||
}
|
}
|
||||||
//send Request
|
//send Request
|
||||||
AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest();
|
AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest();
|
||||||
Collection<AssetOperationRequest> requests = new ArrayList<>();
|
Collection<AssetOperationRequest> requests = new ArrayList<>();
|
||||||
AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart(
|
AssetOperationRequest req = reqBuilder.get();
|
||||||
om_t.getId(),
|
entityId = req.getEntityId();
|
||||||
null,
|
|
||||||
om_t.getBalance(),
|
|
||||||
InOutDirection.out,
|
|
||||||
om_t.getTradingCode(),
|
|
||||||
om_t.getTradingClearingRegistry());
|
|
||||||
requests.add(req);
|
requests.add(req);
|
||||||
|
|
||||||
assetOperationListRequest.setAssetOperationRequests(requests);
|
assetOperationListRequest.setAssetOperationRequests(requests);
|
||||||
log.debug("sending om*t.id={} to gateway", om_t.getId());
|
log.debug("sending {}.id={} to gateway", getName(), entityId);
|
||||||
kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, assetOperationListRequest);
|
kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, assetOperationListRequest);
|
||||||
this.omtId = om_t.getId();
|
this.omtId = req.getEntityId();
|
||||||
boolean gatewayReceived = gatewayCondition.await(GATEWAY_TIMEOUT, TimeUnit.SECONDS);
|
boolean gatewayReceived = gatewayCondition.await(GATEWAY_TIMEOUT, TimeUnit.SECONDS);
|
||||||
if (gatewayReceived) {
|
if (gatewayReceived) {
|
||||||
responseReceived = assetOperationApprovalRequest;
|
responseReceived = assetOperationApprovalRequest;
|
||||||
|
|
@ -103,10 +97,10 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean
|
||||||
}
|
}
|
||||||
if (responseReceived != null) {
|
if (responseReceived != null) {
|
||||||
boolean approved = responseReceived.isApproved();
|
boolean approved = responseReceived.isApproved();
|
||||||
log.debug("om*t.id={} answer from gateway received. approval: {}", om_t.getId(), approved);
|
log.debug("{}.id={} answer from gateway received. approval: {}", getName(), entityId, approved);
|
||||||
return Optional.of(approved);
|
return Optional.of(approved);
|
||||||
} else {
|
} else {
|
||||||
log.error("om*t.id={} answer from gateway not received.", om_t.getId());
|
log.error("{}.id={} answer from gateway not received.", getName(), entityId);
|
||||||
return Optional.empty();
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -122,13 +116,13 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean
|
||||||
gatewayLock.lock();
|
gatewayLock.lock();
|
||||||
try {
|
try {
|
||||||
if (omtId == null) {
|
if (omtId == null) {
|
||||||
log.trace("gateway BaseRequest.id={} currently not waiting for OM*T answer", sysReq.getId());
|
log.trace("gateway BaseRequest.id={} currently not waiting for {} answer", sysReq.getId(), getName());
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
Optional<SingleAssetResponse> assetResponseFound =
|
Optional<SingleAssetResponse> assetResponseFound =
|
||||||
approvals.stream().filter(a -> omtId.equals(a.getEntityId())).findFirst();
|
approvals.stream().filter(a -> omtId.equals(a.getEntityId())).findFirst();
|
||||||
if (assetResponseFound.isPresent()) {
|
if (assetResponseFound.isPresent()) {
|
||||||
log.debug("gateway BaseRequest.id={} OM*T answer found for OM*T.id={}", sysReq.getId(), omtId);
|
log.debug("gateway BaseRequest.id={} answer found for {}.id={}", sysReq.getId(), getName(), omtId);
|
||||||
omtId = null;
|
omtId = null;
|
||||||
assetOperationApprovalRequest = assetResponseFound.get();
|
assetOperationApprovalRequest = assetResponseFound.get();
|
||||||
gatewayCondition.signal();
|
gatewayCondition.signal();
|
||||||
|
|
@ -137,4 +131,12 @@ public class GatewayRequester extends QueueConsumer implements InitializingBean
|
||||||
gatewayLock.unlock();
|
gatewayLock.unlock();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public String getName() {
|
||||||
|
return name != null ? name : "entity";
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setName(String name) {
|
||||||
|
this.name = name;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -12,6 +12,7 @@ import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
import ru.spcex.clearing.error.ClearingError;
|
import ru.spcex.clearing.error.ClearingError;
|
||||||
import ru.spcex.clearing.error.RgsError;
|
import ru.spcex.clearing.error.RgsError;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
|
||||||
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
|
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
|
||||||
import ru.spcex.clearing.service.registry.RegistryManager;
|
import ru.spcex.clearing.service.registry.RegistryManager;
|
||||||
import ru.spcex.clearing.service.schedule.TradingTimeService;
|
import ru.spcex.clearing.service.schedule.TradingTimeService;
|
||||||
|
|
@ -64,6 +65,7 @@ public class InspectionObligations implements ISessionStage {
|
||||||
this.registryManager = registryManager;
|
this.registryManager = registryManager;
|
||||||
this.assets = assets;
|
this.assets = assets;
|
||||||
this.gateway = gateway;
|
this.gateway = gateway;
|
||||||
|
this.gateway.setName("InspectionObligations|OM*T");
|
||||||
this.tradingTimeService = new TradingTimeService(imdgProvider);
|
this.tradingTimeService = new TradingTimeService(imdgProvider);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -133,7 +135,7 @@ public class InspectionObligations implements ISessionStage {
|
||||||
|
|
||||||
Optional<Registry> omt = group.stream().filter(rgs -> RegistryManager.equalsByCode(OM_T, rgs)).findFirst();
|
Optional<Registry> omt = group.stream().filter(rgs -> RegistryManager.equalsByCode(OM_T, rgs)).findFirst();
|
||||||
if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime() && omt.isPresent()) {
|
if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime() && omt.isPresent()) {
|
||||||
Optional<Boolean> gatewayReceived = gateway.gatewayRequestAndWait(omt.get());
|
Optional<Boolean> gatewayReceived = gateway.gatewayRequestAndWait(() -> GatewayRequestCreator.from(omt.get()));
|
||||||
if (gatewayReceived.isEmpty()) {
|
if (gatewayReceived.isEmpty()) {
|
||||||
log.error("Gateway not received response for groupId: {}. ", entry.getKey());
|
log.error("Gateway not received response for groupId: {}. ", entry.getKey());
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -9,6 +9,7 @@ import org.springframework.stereotype.Service;
|
||||||
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
|
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
|
||||||
import ru.clearing.classes.statics.data.registry.Registry;
|
import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.service.integration.GatewayRequestCreator;
|
||||||
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
|
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
|
||||||
import ru.spcex.clearing.service.registry.RegistryManager;
|
import ru.spcex.clearing.service.registry.RegistryManager;
|
||||||
import ru.spcex.clearing.service.schedule.TradingTimeService;
|
import ru.spcex.clearing.service.schedule.TradingTimeService;
|
||||||
|
|
@ -52,6 +53,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage {
|
||||||
this.assets = assets;
|
this.assets = assets;
|
||||||
this.tradingTimeService = tradingTimeService;
|
this.tradingTimeService = tradingTimeService;
|
||||||
this.gateway = gateway;
|
this.gateway = gateway;
|
||||||
|
this.gateway.setName("InspectionObligationsDepositReturn|OM*T");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -66,7 +68,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage {
|
||||||
|
|
||||||
private boolean gatewayIfNeeded(Registry omt) {
|
private boolean gatewayIfNeeded(Registry omt) {
|
||||||
if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime()) {
|
if (SessionType.FINL.equals(sessionType) && tradingTimeService.isTradingTime()) {
|
||||||
Optional<Boolean> gatewayReceived = gateway.gatewayRequestAndWait(omt);
|
Optional<Boolean> gatewayReceived = gateway.gatewayRequestAndWait(() -> GatewayRequestCreator.from(omt));
|
||||||
if (gatewayReceived.isEmpty()) {
|
if (gatewayReceived.isEmpty()) {
|
||||||
log.error("gateway not received response for om*t.id: {}. ", omt.getId());
|
log.error("gateway not received response for om*t.id: {}. ", omt.getId());
|
||||||
} else {
|
} else {
|
||||||
|
|
|
||||||
|
|
@ -45,7 +45,7 @@
|
||||||
<appender-ref ref="CONSOLE"/>
|
<appender-ref ref="CONSOLE"/>
|
||||||
</logger>
|
</logger>
|
||||||
|
|
||||||
<logger name="ru.spcex.clearing.session.stage.impl.GatewayRequester" level="info" additivity="false">
|
<logger name="ru.spcex.clearing.session.stage.impl.GatewayRequester" level="debug" additivity="false">
|
||||||
<appender-ref ref="FILE"/>
|
<appender-ref ref="FILE"/>
|
||||||
<appender-ref ref="CONSOLE"/>
|
<appender-ref ref="CONSOLE"/>
|
||||||
</logger>
|
</logger>
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue