правки валидации sdf10

определение запрос для sdf06/sdf10 - Sdf10Executor и Sdf06Executor
правки запроса к gateway (поле section)
This commit is contained in:
ialbert 2023-08-16 15:54:08 +03:00
parent cef5460726
commit af70ae6392
6 changed files with 66 additions and 23 deletions

View file

@ -4,6 +4,7 @@ import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.AccountBalance;
import ru.clearing.classes.statics.data.account.DepoAccount;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.CompanySymbols;
import ru.clearing.classes.statics.data.company.relation.Relation;
@ -46,6 +47,7 @@ public class ValidationConfig {
Imdg<CompanySymbols> imdgCompanySymbols;
Imdg<Registry> imdgRegistry;
Imdg<Account> imdgAccount;
Imdg<DepoAccount> imdgDepoAccount;
Imdg<AccountBalance> imdgAccountBalance;
Imdg<Security> imdgSecurity;
Imdg<MoneyMarketSecurity> imdgMoneyMarketSecurity;
@ -70,6 +72,7 @@ public class ValidationConfig {
this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class);
this.imdgStatement = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.imdgDepoAccount = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class);
this.imdgProvider = imdgProvider;
}
@ -272,6 +275,11 @@ public class ValidationConfig {
ImdgValidationContext<SDf10> context = new ImdgValidationContext<>();
context.setValidatedObject(sDf10);
context.addImdg(IMDGDistributedNames.Map_Account, imdgAccount);
context.addImdg(IMDGDistributedNames.Map_DepoAccount, imdgDepoAccount);
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
context.addImdg(IMDGDistributedNames.Map_FixedIncomeSecurity , fixedIncomeSecurityImdg);
context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity , imdgMoneyMarketSecurity);
context.addImdg(IMDGDistributedNames.Map_EquitySecurity , equitySecurityImdg);
context.addImdg(IMDGDistributedNames.Map_CompanySymbols, imdgCompanySymbols);
context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
context.addImdg(IMDGDistributedNames.Map_Registry, imdgRegistry);
@ -281,6 +289,7 @@ public class ValidationConfig {
Sdf10ValidationRule.DepoAccountPresent,
CompanyByTradingCodeValidationRule.instance(SDf10::getDepoCode),
Sdf10ValidationRule.CompanyStatus,
Sdf10ValidationRule.TcrPresent,
SecurityBySecurityCodeValidationRule.instance(SDf10::getSecurityCode),
Sdf10ValidationRule.Balance
);

View file

@ -19,6 +19,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandR
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.service.executors.Sdf06Executor;
import ru.spcex.clearing.service.executors.Sdf10Executor;
import ru.spcex.clearing.session.stage.*;
import ru.spcex.clearing.session.stage.impl.BalanceRevise;
import ru.spcex.platform.enumeration.Task;
@ -39,6 +40,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
private final ReturnDepositSession returnDepositSession;
private final SessionManager sessionManager;
private final Sdf06Executor sdf06Executor;
private final Sdf10Executor sdf10Executor;
private final BalanceRevise balanceRevise;
private final Sdf05Sender sdf05Sender;
@ -50,7 +52,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
SecondaryAuctionT0Session secondaryAuctionT0Session,
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
Sdf06Executor sdf06Executor,
BalanceRevise balanceRevise, Sdf05Sender sdf05Sender) {
Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender) {
super(kafkaQueue, kafkaResponseQueue);
this.clearingService = clearingService;
this.registryService = registryService;
@ -63,6 +65,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
this.returnDepositSession = returnDepositSession;
this.sessionManager = sessionManager;
this.sdf06Executor = sdf06Executor;
this.sdf10Executor = sdf10Executor;
this.balanceRevise = balanceRevise;
this.sdf05Sender = sdf05Sender;
}
@ -102,11 +105,14 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
callback(StatementRequest.class)
.setConsumer(sdf06Executor::createStatementAndSendGatewayCommand)
.setConsumer(sdf06Executor::execute)
.forDestination(Consts.STATEMENT_PROCESS_SDF06, callbacks::put);
callback(AssetOperationApprovalRequest.class)
.setConsumer(sdf06Executor::processGatewayResponse)
.setConsumer(req -> {
sdf06Executor.processGatewayResponse(req);
sdf10Executor.processGatewayResponse(req);
})
.forDestination(Consts.ASSET_OPERATION_APPROVAL, callbacks::put);
callback(STradesImportedRequest.class)

View file

@ -31,6 +31,7 @@ import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.number.BigDecimalUtil;
@ -102,7 +103,7 @@ public class Sdf06Executor {
this.filenameObtainer = filenameObtainer;
}
public void createStatementAndSendGatewayCommand(BaseRequest<StatementRequest> systemRequest) {
public void execute(BaseRequest<StatementRequest> systemRequest) {
Imdg<SDf06> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf06, SDf06.class);
Long groupId = systemRequest.getRequestPayload().getGroupId();
Collection<SDf06> sdfs = sdfImdg.getCollectionObjectsByFieldValues(Map.of(
@ -209,15 +210,23 @@ public class Sdf06Executor {
//может прийти неограниченно позже 06го, после того как прошел клиринг например...
Instant now = Instant.now();
String fileName = null;
boolean requestIsIntendedForSdf06 = false;
for (SingleAssetResponse gatewayMsg : req.getRequestPayload().getApprovals()) {
//получаем запрос для текущей группы sdf06
//находим группу
Long statementId = gatewayMsg.getStatementId();
Statement stmt = statementImdg.getSingleObjectByID(statementId);
ImdgPredicateBuilder pb = statementImdg.predicateBuilder();
Statement stmt = statementImdg.getFirstObjectByPredicate(
pb.and(
pb.equals("id", statementId),
pb.equals("inOutSDfType", InOutSDfType.type6.getKey())
)
);
if (stmt == null) {
log.error("Statement.id {} not found", statementId);
log.debug("Statement.id {} with type {} not found", statementId, InOutSDfType.type6.getKey());
return;
}
requestIsIntendedForSdf06 = true;
Long sdf06Id = stmt.getInSDfId();
SDf06 sdf06 = sdf06Imdg.getSingleObjectByID(sdf06Id);
@ -256,9 +265,13 @@ public class Sdf06Executor {
statementImdg.update(stmt);
}
}
sendToExporter(sdf07GroupId, fileName);
log.debug("All gateway responses received for SDF06 groupId {}. SDF07 groupId {}", sdf06GroupId, sdf07GroupId);
clearContext();
if (requestIsIntendedForSdf06) {
sendToExporter(sdf07GroupId, fileName);
log.debug("All gateway responses received for SDF06 groupId {}. SDF07 groupId {}", sdf06GroupId, sdf07GroupId);
clearContext();
} else {
log.debug("gateway response was not for SDF06 executor");
}
}
private void clearContext() {

View file

@ -31,6 +31,7 @@ import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.number.BigDecimalUtil;
@ -155,7 +156,7 @@ public class Sdf10Executor {
processedApproved(stmt, sDf10, now, sdf11GroupId);
sdf11WasCreated = true;
} else {
requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode()));
requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode(), sDf10.getSecurityCode()));
}
}
if (requests.size() > 0) {
@ -189,15 +190,23 @@ public class Sdf10Executor {
public void processGatewayResponse(BaseRequest<AssetOperationApprovalRequest> req) {
Instant now = Instant.now();
boolean requestIsIntendedForSdf10 = false;
for (SingleAssetResponse gatewayMsg : req.getRequestPayload().getApprovals()) {
//получаем запрос для текущей группы sdf10
//находим группу
Long statementId = gatewayMsg.getStatementId();
Statement stmt = statementImdg.getSingleObjectByID(statementId);
ImdgPredicateBuilder pb = statementImdg.predicateBuilder();
Statement stmt = statementImdg.getFirstObjectByPredicate(
pb.and(
pb.equals("id", statementId),
pb.equals("inOutSDfType", InOutSDfType.type10.getKey())
)
);
if (stmt == null) {
log.error("Statement.id {} not found", statementId);
log.debug("Statement.id {} with type {} not found", statementId, InOutSDfType.type6.getKey());
return;
}
requestIsIntendedForSdf10 = true;
Long sdf10Id = stmt.getInSDfId();
SDf10 sdf10 = sdf10Imdg.getSingleObjectByID(sdf10Id);
@ -233,9 +242,13 @@ public class Sdf10Executor {
statementImdg.update(stmt);
}
}
sendToExporter(sdf11GroupId);
log.debug("All gateway responses received for SDF10 groupId {}. SDF11 groupId {}", sdf10GroupId, sdf11GroupId);
clearContext();
if (requestIsIntendedForSdf10) {
sendToExporter(sdf11GroupId);
log.debug("All gateway responses received for SDF10 groupId {}. SDF11 groupId {}", sdf10GroupId, sdf11GroupId);
clearContext();
} else {
log.debug("gateway response was not for SDF10 executor");
}
}
private void clearContext() {
@ -280,17 +293,18 @@ public class Sdf10Executor {
stmt.setAmount(amount.abs());
stmt.setOperationStatus(OperationStatus.Pending.getKey());
stmt.setInSDfId(sdf10.getId());
stmt.setInOutSDfType(InOutSDfType.type6.getKey());
stmt.setInOutSDfType(InOutSDfType.type10.getKey());
stmt.setClearingDate(LocalDate.now());
stmt.setCreated(Instant.now());
return stmt;
}
private AssetOperationRequest requestFromStatement(Statement stmt, String tradingCode, String tcrCode) {
private AssetOperationRequest requestFromStatement(Statement stmt, String tradingCode, String tcrCode, String securityCode) {
AssetOperationRequest req = new AssetOperationRequest();
req.setStatementId(stmt.getId());
req.setAmount(stmt.getAmount());
req.setSecuritySymbol(CurrencyCode.RUB.getKey()); //fixme retrieve security symbol from validator
// req.setAmount(stmt.getAmount());
req.setQuantity(stmt.getAmount());
req.setSecuritySymbol(securityCode); //fixme retrieve security symbol from validator
req.setTradingCode(tradingCode);
req.setCode(tcrCode);
req.setDirection(stmt.getInOutDirection());

View file

@ -36,7 +36,7 @@ public enum Sdf10ValidationRule implements IValidationRule<ImdgValidationContext
if (TextUtil.isEmpty(sdf10.getDepoCode())) {
return of(ClearingError.DepoAccNotFound, sdf10.getDepoCode());
}
Imdg<Account> accountImdg = context.obtainMap(IMDGDistributedNames.Map_DepoAccount, Account.class);
Imdg<Account> accountImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class);
Imdg<DepoAccount> depoAccountImdg = context.obtainMap(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class);
Collection<Account> accs = accountImdg.getCollectionObjectsBySQL("account = '%s'".formatted(sdf10.getDepoCode()));
if (accs.size() != 1) {
@ -68,7 +68,7 @@ public enum Sdf10ValidationRule implements IValidationRule<ImdgValidationContext
TcrPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf10> context) {
DepoAccount account = context.getStoredObject(ValidationStored.Sdf10Account);
Account account = context.getStoredObject(ValidationStored.Sdf10Account);
if (account == null) {
return of(ClearingError.TCRegistryNotFound, "<no account to search for>");
}
@ -92,7 +92,7 @@ public enum Sdf10ValidationRule implements IValidationRule<ImdgValidationContext
// }
// }
TradingClearingRegistry tcr = tcrImdg.getFirstObjectByFieldValues(
Map.of("moneyAccountId", account.getId())
Map.of("depoAccountId", account.getId())
);
if (tcr == null) {
return of(ClearingError.TCRegistryNotFound, "<moneyAccountId = " + account.getId() + ">");
@ -114,6 +114,7 @@ public enum Sdf10ValidationRule implements IValidationRule<ImdgValidationContext
ImdgPredicateBuilder prdctBldr = registryImdg.predicateBuilder();
Function<RegistryTradingParams, ImdgPredicate> prdctByRegistryCode = code -> prdctBldr.and(
prdctBldr.equals("account", account.getAccount()),
prdctBldr.equals("securitySymbol", sdf10.getSecurityCode()),
prdctBldr.equals("companyId", company.getId()),
prdctBldr.sql(RegistryCodeSqlBuilder.getInstance(code).build()));
Registry rgsAsf = registryImdg.getSingleObjectByPredicate(prdctByRegistryCode.apply(RegistryTradingParams.AS_F));

View file

@ -183,7 +183,7 @@ public class GatewayService extends QueueConsumer implements InitializingBean {
}
OutboundRequest outboundRequest = OutboundRequestBuilder.builder()
.section(Section.MKR.getKey())
.section(assetOperation.getAmount() != null ? Section.MKR.getKey() : Section.FOND.getKey())
.type(OutboundRequestType.ASSET_OPERATION.getKey())
.content(content).build();