statement services fixes

This commit is contained in:
ialbert 2023-08-09 21:41:22 +03:00
parent 3f3c59cc00
commit fbfe589ea3
19 changed files with 557 additions and 74 deletions

View file

@ -199,9 +199,9 @@ public class ValidationConfig {
context.addImdg(IMDGDistributedNames.Map_EquitySecurity, equitySecurityImdg);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
Sdf21ValidationRule.CompanyPresent,
CompanyByDepoCodeValidationRule.instance(SDf21::getDepoCodeCl),
Sdf21ValidationRule.AccountPresent,
Sdf21ValidationRule.InstrumentPresent
SecurityBySecurityCodeValidationRule.instance(SDf21::getSecurityCode)
);
};
}
@ -261,6 +261,26 @@ public class ValidationConfig {
};
}
@Bean("sdf10Validator")
public Function<SDf10, IValidator> sdf10Validator() {
return sDf10 -> {
ImdgValidationContext<SDf10> context = new ImdgValidationContext<>();
context.setValidatedObject(sDf10);
context.addImdg(IMDGDistributedNames.Map_Account, imdgAccount);
context.addImdg(IMDGDistributedNames.Map_CompanySymbols, imdgCompanySymbols);
context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
context.addImdg(IMDGDistributedNames.Map_Registry, imdgRegistry);
context.addImdg(IMDGDistributedNames.Map_Session, imdgSession);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
Sdf10ValidationRule.DepoAccountPresent,
CompanyByDepoCodeValidationRule.instance(SDf10::getDepoCode),
Sdf10ValidationRule.CompanyStatus,
SecurityBySecurityCodeValidationRule.instance(SDf10::getSecurityCode)
);
};
}
@Bean("returnDepositValidator")
public Function<RegistryReturnDepositRequest, IValidator> returnDepositValidator() {
return returnDepositRequest -> {

View file

@ -4,6 +4,7 @@ import ru.spcex.platform.utils.enumeration.IErrorEnumId;
public enum ClearingError implements IErrorEnumId {
GeneralError(5400L),
DepoAccNotFound(5011L),
RecordNotFoundInDictionary(5423L),
IncorrectValue(5404L),
RecordNotFound(5406L),

View file

@ -71,11 +71,11 @@ public class StatementService extends QueueConsumer implements InitializingBean
.setConsumer(systemRequest -> {
StatementRequest payload = systemRequest.getRequestPayload();
if (payload.getTable() == null
|| !Arrays.asList(SdfTable.SDF_01, SdfTable.SDF_04, SdfTable.SDF_57, SdfTable.SDF_08, SdfTable.SDF_13, SdfTable.SDF_21).contains(payload.getTable())) {
|| !Arrays.asList(SdfTable.SDF_01, SdfTable.SDF_04, SdfTable.SDF_57, SdfTable.SDF_08, SdfTable.SDF_13, SdfTable.SDF_21, SdfTable.SDF_10).contains(payload.getTable())) {
log.debug("StatementService: skipping table {}", payload.getTable());
return;
}
if (Arrays.asList(SdfTable.SDF_08, SdfTable.SDF_21).contains(payload.getTable())) {
if (Arrays.asList(SdfTable.SDF_08, SdfTable.SDF_21, SdfTable.SDF_10).contains(payload.getTable())) {
stmtSrvV2.processReq(systemRequest);
} else if (payload.isContinueSdf()) {
processAccountAnswer(systemRequest);
@ -295,7 +295,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
AbstractExecutor service = executorsMap.get(SdfTable.SDF_01);
Result res = service.execute(sdfGroup, statementRequest);
if (res.getAccountRequests().size() != 0) {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests(), res.getGenerationId()));
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests(), res.getChildGenerationId()));
return res;
} else if (service.isNeedToSendCommand()) {
service.sendCommand(kafkaSender, res);

View file

@ -12,6 +12,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf0
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.executors.Sdf08Executor;
import ru.spcex.clearing.service.executors.Sdf10Executor;
import ru.spcex.clearing.service.executors.Sdf21Executor;
import ru.spcex.clearing.service.model.Result;
import ru.spcex.platform.enumeration.SdfTable;
@ -19,6 +20,7 @@ import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.util.*;
import java.util.function.Predicate;
import java.util.stream.Collectors;
@Component
@ -29,12 +31,14 @@ public class StatementServiceV2 {
private final List<StatementRequest> statementRequests = new LinkedList<>();
private final Sdf08Executor sdf08Executor;
private final Sdf21Executor sdf21Executor;
private final Sdf10Executor sdf10Executor;
public StatementServiceV2(ImdgProvider imdgProvider, KafkaSender kafkaSender, Sdf08Executor sdf08Executor, Sdf21Executor sdf21Executor) {
public StatementServiceV2(ImdgProvider imdgProvider, KafkaSender kafkaSender, Sdf08Executor sdf08Executor, Sdf21Executor sdf21Executor, Sdf10Executor sdf10Executor) {
this.imdgProvider = imdgProvider;
this.kafkaSender = kafkaSender;
this.sdf08Executor = sdf08Executor;
this.sdf21Executor = sdf21Executor;
this.sdf10Executor = sdf10Executor;
}
public void processReq(BaseRequest<StatementRequest> systemRequest) {
@ -45,12 +49,18 @@ public class StatementServiceV2 {
if (sdfGroup.isEmpty()) {
log.debug("table {} is not paired with any other SDF", table);
//здесь switch case одиночные методы
switch (table) {
case SDF_10 -> {
sdf10Executor.execute(systemRequest);
}
default -> log.error("unknown table {}", table);
}
return;
}
Collection<StatementRequest> fullGroup = getFullGroup(systemRequest.getRequestPayload(), sdfGroup.get());
statementRequests.add(systemRequest.getRequestPayload());
if (fullGroup.isEmpty()) {
log.debug("didn't find full sdf set for stmtReq.SdfTable={}, groupId={}. Adding to cache", table, groupId);
statementRequests.add(systemRequest.getRequestPayload());
return;
}
log.debug("find full set of SDF requests: {}", fullGroup.stream()
@ -78,16 +88,28 @@ public class StatementServiceV2 {
*/
private void processSdf08And21(StatementRequest sdf08, StatementRequest sdf21) {
Result sdf08Res = processSdf08(sdf08);
removeFirstWithSameTableAndGroupId(sdf08);
if (sdf08Res.getAccountRequests().size() > 0) {
statementRequests.remove(sdf08);
log.info("sdf08 execution wasn't complete, waiting for an answer from account-service");
return;
}
//затем sdf21
processSdf21(sdf21);
removeFirstWithSameTableAndGroupId(sdf21);
//fixme ревизия для бумаг reviser.doRevise(pair.getFirst().getGroupId());
log.info("pair sdf08/sdf21 processed successfully");
statementRequests.removeAll(Arrays.asList(sdf08, sdf21));
}
private void removeFirstWithSameTableAndGroupId(StatementRequest statementRequest) {
Iterator<StatementRequest> iterator = statementRequests.iterator();
Predicate<StatementRequest> eqs = r -> Objects.equals(r.getTable(), statementRequest.getTable()) && Objects.equals(r.getGroupId(), statementRequest.getGroupId());
while (iterator.hasNext()) {
StatementRequest next = iterator.next();
if (eqs.test(next)) {
iterator.remove();
return;
}
}
}
@ -108,7 +130,7 @@ public class StatementServiceV2 {
if (res.getAccountRequests().size() != 0) {
AccountSdf01Request createAccsReq = StatementService.createAccountsRequest(statementRequest.getGroupId(),
res.getAccountRequests(),
res.getGenerationId());
res.getChildGenerationId());
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF08, createAccsReq);
} else if (sdf08Executor.isNeedToSendCommand()) {
sdf08Executor.sendCommand(kafkaSender, res);

View file

@ -97,7 +97,7 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
@Override
public void sendCommand(KafkaSender kafkaSender, Result result) {
ExportToFileRequest exportRequest = new ExportToFileRequest();
exportRequest.setSdfGroupId(result.getGenerationId());
exportRequest.setSdfGroupId(result.getChildGenerationId());
exportRequest.setNameOfTable(exportTableName());
kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest);
}
@ -105,7 +105,7 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
public Result execute(Collection<SDf01> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = statementRequest.getChildGenerationId() != null ? statementRequest.getChildGenerationId() : imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
result.setChildGenerationId(generationIdForGroup);
// boolean reviseFailed = false;
log.info("SDF01 execution: sdf01 number={}, groupId={}", sdf.size(), sdf.stream().findFirst().map(SDf01::getGenerationId).orElse(null));
for (SDf01 sdf01 : sdf) {

View file

@ -70,7 +70,7 @@ public class Sdf04Executor extends AbstractExecutor<SDf04> {
public Result execute(Collection<SDf04> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
result.setChildGenerationId(generationIdForGroup);
for (SDf04 sdf04 : sdf) {
//обычно мы ищем по группу sdf04, здесь как будто всегда только одна запись, todo нужно прочекать этот момент
log.debug("Process sdf04 record; sdf04.id: {}", sdf04.getId());

View file

@ -96,7 +96,7 @@ public class Sdf08Executor extends AbstractExecutor<SDf08> {
public void sendCommand(KafkaSender kafkaSender, Result result) {
SwtExporterRequest swtReq = new SwtExporterRequest();
swtReq.setType(SdfTable.SDF_09.getKey());
swtReq.setGroupId(result.getGenerationId());
swtReq.setGroupId(result.getChildGenerationId());
kafkaSender.sendRequestToQueue(Consts.SWT_EXPORTER, swtReq);
}
@ -123,8 +123,9 @@ public class Sdf08Executor extends AbstractExecutor<SDf08> {
public Result execute(Collection<SDf08> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
Long generationIdForGroup = statementRequest.getChildGenerationId() != null ?
statementRequest.getChildGenerationId() :imdgProvider.getImdgIdGenerator().nextId();
result.setChildGenerationId(generationIdForGroup);
for (SDf08 sdf08 : sdf) {
Statement stmt = null;
try {

View file

@ -0,0 +1,322 @@
package ru.spcex.clearing.service.executors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
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.DepoAccount;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.classes.statics.data.sdf.SDf06;
import ru.clearing.classes.statics.data.sdf.SDf07;
import ru.clearing.classes.statics.data.sdf.SDf10;
import ru.clearing.classes.statics.data.sdf.SDf11;
import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
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.AssetOperationRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
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.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Collection;
import java.util.Map;
import java.util.function.Function;
@Service
public class Sdf10Executor {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<Statement> statementImdg;
private final ImdgProvider imdgProvider;
private final Imdg<Registry> registryImdg;
private final Imdg<FixedIncomeSecurity> fixedIncomeSecurityImdg;
private final Imdg<EquitySecurity> equitySecurityImdg;
private final Imdg<DepoAccount> depoAccountImdg;
private final Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
private final Imdg<Company> companyImdg;
private final Imdg<SDf06> sdf06Imdg;
private final Imdg<SDf07> sdf07Imdg;
private final ImdgId idGenerator;
private final IMessageResolver messageResolver;
private final Function<SDf06, IValidator> sDf06Validator;
private final KafkaSender kafkaSender;
private final static BigDecimal successResult = BigDecimal.ZERO;
//ошибка проверок
private final static BigDecimal errorResult2 = new BigDecimal("3");
//ошибка проверок баланса или наличия сессии - в этом сценарии session создается
private final static BigDecimal errorResult1 = new BigDecimal("1");
//отказ от gateway
private final static BigDecimal errorResult3 = new BigDecimal("3");
private Long sdf10GroupId;
private Long sdf07GroupId;
@Autowired
public Sdf10Executor(ImdgProvider imdgProvider,
IMessageResolver messageResolver,
@Qualifier("sdf06ValidatorNew") Function<SDf06, IValidator> sDf06Validator,
KafkaSender kafkaSender) {
this.imdgProvider = imdgProvider;
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class);
this.depoAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class);
this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.messageResolver = messageResolver;
this.sDf06Validator = sDf06Validator;
this.kafkaSender = kafkaSender;
this.sdf06Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf06, SDf06.class);
this.sdf07Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf07, SDf07.class);
}
public void execute(BaseRequest<StatementRequest> systemRequest) {
Imdg<SDf10> sdf10Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf10, SDf10.class);
Imdg<SDf11> sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class);
Long groupId = systemRequest.getRequestPayload().getGroupId();
Collection<SDf10> sdfs = sdf10Imdg.getCollectionObjectsByFieldValues(Map.of(
"generationId", groupId
));
Instant now = Instant.now();
Long sdf11GroupId = idGenerator.nextId();
log.info("processing SDF10 groupId: {} sdf10s: {}", groupId, sdfs.size());
for (SDf10 sdf10 : sdfs) {
SDf11 sdf11 = createSdf11(sdf10, now);
sdf11.setGenerationId(sdf11GroupId);
sdf11.setResult("OK");
sdf11Imdg.insert(sdf11);
log.debug("sdf10.id={} created sdf11.id={}", sdf10.getId(), sdf11.getId());
}
sendToExporter(sdf11GroupId);
// if (sdf10GroupId != null) {
// log.warn("currently awaiting gateway response for groupId: {}. skipping groupId {}", sdf10GroupId, groupId);
// sdfs.forEach(sdf10 -> {
// SDf11 errorSdf11 = createSdf11(sdf10, now, errorResult2);
// errorSdf11.setGenerationId(sdf11GroupId);
// sdf07Imdg.insert(errorSdf11);
// });
// sendToExporter(sdf11GroupId);
// return;
// }
// Collection<AssetOperationRequest> requests = new ArrayList<>();
// boolean sdf07WasCreated = false;
// for (SDf06 sDf06 : sdfs) {
// IValidator validator = sDf06Validator.apply(sDf06);
// Optional<EnumMessage> err = validator.tillFirstError();
// if (err.isPresent()) {
// log.debug("sdf06.id={} validation error: {}", sDf06.getId(), messageResolver.resolve(err.get()));
// SDf07 errorSdf07;
// BigDecimal errorResult;
// Statement stmt = null;
// if (needToCreateStatement(err.get())) {
// errorResult = errorResult1;
// stmt = createStatementBySdf06(sDf06,
// ((Company) validator.getStored(ValidationStored.Sdf06Company)).getId(),
// validator.getStored(ValidationStored.Sdf06Account));
// } else {
// errorResult = errorResult2;
// }
// errorSdf07 = createSdf11(sDf06, now, errorResult);
// errorSdf07.setGenerationId(sdf11GroupId);
// sdf07Imdg.insert(errorSdf07);
// sdf07WasCreated = true;
// if (stmt != null) {
// stmt.setOutSDfId(errorSdf07.getId());
// statementImdg.insert(stmt);
// log.debug("Statement created: {}", stmt.getId());
// }
// continue;
// }
// Company company = validator.getStored(ValidationStored.Sdf06Company);
// Account account = validator.getStored(ValidationStored.Sdf06Account);
// TradingClearingRegistry tcr = tradingClearingRegistryImdg.getFirstObjectByFieldValues(
// Map.of("moneyAccountId", account.getId())
// );
// //проверка существует ли statement пока убрал
// Statement stmt = createStatementBySdf06(sDf06, company.getId(), account);
// statementImdg.insert(stmt);
// log.debug("Statement created: {}", stmt.getId());
// if (InOutDirection.in.getKey().equals(stmt.getInOutDirection())) {
// processedApproved(stmt, sDf06, now, sdf11GroupId);
// sdf07WasCreated = true;
// } else {
// requests.add(requestFromStatement(stmt, company.getTradingCode(), tcr.getCode()));
// }
// }
// if (requests.size() > 0) {
// this.sdf10GroupId = groupId;
// this.sdf07GroupId = sdf11GroupId;
// AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest();
// assetOperationListRequest.setAssetOperationRequests(requests);
// kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, assetOperationListRequest);
// } else if (sdf07WasCreated) {
// sendToExporter(sdf11GroupId);
// }
}
// private void processedApproved(Statement stmt, SDf06 sDf06, Instant time) {
// processedApproved(stmt, sDf06, time, sdf07GroupId);
// }
//
// private void processedApproved(Statement stmt, SDf06 sDf06, Instant time, Long sdf07GenerationId) {
// log.trace("statement.id={}, sdf07.id={}, sdf06.id={} executed",
// stmt.getId(),
// sDf06.getGenerationId(),
// sDf06.getId());
// SDf07 sdf07 = createSdf11(sDf06, time, successResult);
// sdf07.setGenerationId(sdf07GenerationId);
// sdf07Imdg.insert(sdf07);
// stmt.setOperationStatus(OperationStatus.Executed.getKey());
// stmt.setUpdated(time);
// stmt.setOutSDfId(sdf07.getId());
// statementImdg.update(stmt);
// }
// private boolean needToCreateStatement(EnumMessage err) {
// return ClearingError.ActiveSessionIsPresent.equals(err.getSubject()) || ClearingError.BalanceInsufficient.equals(err.getSubject());
// }
//
// public void processGatewayResponse(BaseRequest<AssetOperationApprovalRequest> req) {
// Instant now = Instant.now();
// for (SingleAssetResponse gatewayMsg : req.getRequestPayload().getApprovals()) {
// //получаем запрос для текущей группы sdf06
// //находим группу
// Long statementId = gatewayMsg.getStatementId();
// Statement stmt = statementImdg.getSingleObjectByID(statementId);
// if (stmt == null) {
// log.error("Statement.id {} not found", statementId);
// return;
// }
//
// Long sdf06Id = stmt.getInSDfId();
// SDf06 sdf06 = sdf06Imdg.getSingleObjectByID(sdf06Id);
// if (sdf06 == null) {
// log.error("Sdf06.id {} not found by statement.id {}", sdf06Id, statementId);
// return;
// }
//
// Long groupId = sdf06.getGenerationId();
// //сверяем группу SDF06 пришедшего запроса с ожидаемой
// if (sdf10GroupId == null || !sdf10GroupId.equals(groupId)) {
// log.error("do not currently waiting for gateway response for statement.id {} sdf06 groupId {}; waiting for {}",
// statementId,
// groupId,
// sdf06Id);
// return;
// }
//
// Instant updatedTime = Instant.now();
// if (gatewayMsg.isApproved()) {
// processedApproved(stmt, sdf06, updatedTime);
// } else {
// log.trace("statement.id={}, sdf07.id={}, sdf06.id={} rejected (by gateway answer)",
// statementId,
// sdf06.getGenerationId(),
// sdf06.getId());
// SDf07 sdf07 = createSdf11(sdf06, now, errorResult3);
// sdf07.setGenerationId(sdf07GroupId);
// sdf07Imdg.insert(sdf07);
// stmt.setOperationStatus(OperationStatus.Rejected.getKey());
// stmt.setUpdated(updatedTime);
// stmt.setOutSDfId(sdf07.getId());
// statementImdg.update(stmt);
// }
// }
// sendToExporter(sdf07GroupId);
// log.debug("All gateway responses received for SDF06 groupId {}. SDF07 groupId {}", sdf10GroupId, sdf07GroupId);
// clearContext();
// }
//
// private void clearContext() {
// this.sdf10GroupId = null;
// this.sdf07GroupId = null;
// }
private void sendToExporter(Long generationId) {
SwtExporterRequest swtReq = new SwtExporterRequest();
swtReq.setType(SdfTable.SDF_11.getKey());
swtReq.setGroupId(generationId);
kafkaSender.sendRequestToQueue(Consts.SWT_EXPORTER, swtReq);
}
private SDf11 createSdf11(SDf10 sdf10, Instant time) {
// SDf11 sDf07 = new SDf11();
// sDf07.setInSDfId(sdf06.getId());
// sDf07.setGenerationTime(time);
// sDf07.setAccount(sdf06.getAccount());
// sDf07.setSum(sdf06.getSum());
// sDf07.setMarket(sdf06.getMarket());
// sDf07.setType(sdf06.getType());
// sDf07.setDeal(sdf06.getDeal());
// sDf07.setClientN(sdf06.getClientN());
// sDf07.setInn(sdf06.getInn());
// sDf07.setBic(sdf06.getBic());
// sDf07.setNumber(sdf06.getNumber());
// sDf07.setSpec(sdf06.getSpec());
// sDf07.setResult(result);
// return sDf07;
SDf11 sdf11 = new SDf11();
sdf11.setOutDocument(sdf10.getOutDocument());
sdf11.setDepoCode(sdf10.getDepoCode());
sdf11.setQuantity(sdf10.getQuantity());
sdf11.setSecurityCode(sdf10.getSecurityCode());
sdf11.setClientName(sdf10.getClientName());
sdf11.setGenerationTime(time);
// sdf11.setGenerationId(sdf10.getGenerationId());
return sdf11;
}
private Statement createStatementBySdf06(SDf06 sdf06, Long companyId, Account account) {
Statement stmt = new Statement();
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());
if (amount.compareTo(BigDecimal.ZERO) >= 0) {
stmt.setInOutDirection(InOutDirection.in.getKey());
} else {
stmt.setInOutDirection(InOutDirection.out.getKey());
}
stmt.setAmount(amount.abs());
stmt.setOperationStatus(OperationStatus.Pending.getKey());
stmt.setInSDfId(sdf06.getId());
stmt.setInOutSDfType(InOutSDfType.type6.getKey());
stmt.setClearingDate(LocalDate.now());
stmt.setCreated(Instant.now());
return stmt;
}
private AssetOperationRequest requestFromStatement(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

@ -61,7 +61,7 @@ public class Sdf13Executor extends AbstractExecutor<SDf13> {
public Result execute(Collection<SDf13> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
result.setChildGenerationId(generationIdForGroup);
// for (SDf13 sdf13 : sdf) {
// //обычно мы ищем по группу sdf04, здесь как будто всегда только одна запись, todo нужно прочекать этот момент
// SDf12 sDf12 = selectSdf12bySdf13(sdf13);

View file

@ -131,7 +131,7 @@ public class Sdf21Executor extends AbstractExecutor<SDf21> {
public Result execute(Collection<SDf21> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
result.setChildGenerationId(generationIdForGroup);
log.info("SDF21 execution: sdf21 number={}, groupId={}", sdf.size(), sdf.stream().findFirst().map(SDf21::getGenerationId).orElse(null));
for (SDf21 sdf21 : sdf) {
IValidator validator = sDf21Validator.apply(sdf21);
@ -140,7 +140,7 @@ public class Sdf21Executor extends AbstractExecutor<SDf21> {
log.error("error while validating sdf21.id={} - {}", sdf21.getId(), messageResolver.resolve(error.get()));
continue;
}
Company company = validator.getStored(ValidationStored.Sdf21Company);
Company company = validator.getStored(ValidationStored.CompanyByDepoCode);
Account account = validator.getStored(ValidationStored.Sdf21Account);
log.debug("company.id={}, account.id={}", company.getId(), account.getId());
//адресат CREDIT / владелец DEBIT

View file

@ -135,7 +135,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
public Result execute(Collection<SDf57> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
result.setChildGenerationId(generationIdForGroup);
log.info("SDF57 execution: sdf57 number={}, groupId={}", sdf.size(), sdf.stream().findFirst().map(SDf57::getGenerationId).orElse(null));
for (SDf57 sdf57 : sdf) {
IValidator validator = sDf57Validator.apply(sdf57);

View file

@ -76,7 +76,7 @@ public class SdfLegacyExecutor extends AbstractExecutor<SDf01> {
@Override
public void sendCommand(KafkaSender kafkaSender, Result result) {
ExportToFileRequest exportRequest = new ExportToFileRequest();
exportRequest.setSdfGroupId(result.getGenerationId());
exportRequest.setSdfGroupId(result.getChildGenerationId());
exportRequest.setNameOfTable(exportTableName());
kafkaSender.sendRequestToQueue(Consts.EXPORT_PROCESS, exportRequest);
}
@ -84,7 +84,7 @@ public class SdfLegacyExecutor extends AbstractExecutor<SDf01> {
public Result execute(Collection<SDf01> sdf, StatementRequest statementRequest) {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
result.setChildGenerationId(generationIdForGroup);
for (SDf01 sdf01 : sdf) {
IValidator validator = sDf01Validator.apply(sdf01);
Optional<EnumMessage> error = validator.tillFirstError();

View file

@ -7,7 +7,7 @@ import java.util.List;
public class Result {
private List<AccountSdfRequestPart> accountRequests = new ArrayList<>();
private Long generationId;
private Long childGenerationId;
public List<AccountSdfRequestPart> getAccountRequests() {
return accountRequests;
@ -17,11 +17,11 @@ public class Result {
this.accountRequests = accountRequests;
}
public Long getGenerationId() {
return generationId;
public Long getChildGenerationId() {
return childGenerationId;
}
public void setGenerationId(Long generationId) {
this.generationId = generationId;
public void setChildGenerationId(Long generationId) {
this.childGenerationId = generationId;
}
}

View file

@ -0,0 +1,47 @@
package ru.spcex.clearing.service.validation;
import ru.clearing.classes.statics.data.company.Company;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Collection;
import java.util.Objects;
import java.util.Optional;
import java.util.function.Function;
public class CompanyByDepoCodeValidationRule<T> implements IValidationRule<ImdgValidationContext<T>> {
private final Function<T, String> depoCodeExtractor;
private CompanyByDepoCodeValidationRule(Function<T, String> depoCodeExtractor) {
Objects.requireNonNull(depoCodeExtractor, "cannot create CompanyByDepoCodeValidationRule: depoCodeExtractor is null");
this.depoCodeExtractor = depoCodeExtractor;
}
public static <C> CompanyByDepoCodeValidationRule<C> instance(Function<C, String> depoCodeExtractor) {
return new CompanyByDepoCodeValidationRule<>(depoCodeExtractor);
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<T> context) {
T validatedObject = context.getValidatedObject();
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
String depoCode = depoCodeExtractor.apply(validatedObject);
if (TextUtil.isEmpty(depoCode) || depoCode.length() < 4) {
return of(ClearingError.CompanyNotFoundB, depoCode);
}
String tradingCode = depoCode.substring(1, 4);
tradingCode = tradingCode.replaceFirst("^0+(?!$)", "");
Collection<Company> companies = companyImdg.getCollectionObjectsBySQL("tradingCode = '%s'".formatted(tradingCode));
if (companies.size() != 1) {
return of(ClearingError.CompanyNotFoundB, tradingCode);
}
context.storeObject(ValidationStored.CompanyByDepoCode, companies.iterator().next());
return empty();
}
}

View file

@ -0,0 +1,62 @@
package ru.spcex.clearing.service.validation;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.DepoAccount;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.sdf.SDf10;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Collection;
import java.util.Optional;
public enum Sdf10ValidationRule implements IValidationRule<ImdgValidationContext<SDf10>> {
DepoAccountPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf10> context) {
SDf10 sdf10 = context.getValidatedObject();
if (TextUtil.isEmpty(sdf10.getDepoCode())) {
return of(ClearingError.DepoAccNotFound, sdf10.getDepoCode());
}
Imdg<Account> accountImdg = context.obtainMap(IMDGDistributedNames.Map_DepoAccount, 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) {
return of(ClearingError.DepoAccNotFound, sdf10.getDepoCode());
}
Account acc = accs.iterator().next();
DepoAccount depoAccount = depoAccountImdg.getFirstObjectBySQL("accountId = %d".formatted(acc.getId()));
if (depoAccount == null) {
return of(ClearingError.DepoAccNotFound, sdf10.getDepoCode());
}
return empty();
}
},
CompanyStatus() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf10> context) {
Company company = context.getStoredObject(ValidationStored.CompanyByDepoCode);
if (company == null) {
return of(ClearingError.CompanyNotActive, "null");
}
if (!WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
return of(ClearingError.CompanyNotActive, company.getTradingCode());
}
return empty();
}
},
;
private final static Logger log = LoggerFactory.getLogger(Sdf10ValidationRule.class);
@Override
public String ruleName() {
return "Sdf10ValidationRule." + name();
}
}

View file

@ -1,14 +1,11 @@
package ru.spcex.clearing.service.validation;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.sdf.SDf21;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.enumeration.AccountType;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.text.TextUtil;
@ -19,26 +16,6 @@ import java.util.Optional;
public enum Sdf21ValidationRule implements IValidationRule<ImdgValidationContext<SDf21>> {
CompanyPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf21> context) {
SDf21 sdf21 = context.getValidatedObject();
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
String depoCodeCl = sdf21.getDepoCodeCl();
//ClearingError.CompanyNotFoundB
if (TextUtil.isEmpty(depoCodeCl) || depoCodeCl.length() < 4) {
return of(ClearingError.CompanyNotFoundB, depoCodeCl);
}
String tradingCode = depoCodeCl.substring(1, 4);
tradingCode = tradingCode.replaceFirst("^0+(?!$)", "");
Collection<Company> companies = companyImdg.getCollectionObjectsBySQL("tradingCode = '%s'".formatted(tradingCode));
if (companies.size() != 1) {
return of(ClearingError.CompanyNotFoundB, depoCodeCl);
}
context.storeObject(ValidationStored.Sdf21Company, companies.iterator().next());
return empty();
}
},
AccountPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf21> context) {
@ -56,28 +33,7 @@ public enum Sdf21ValidationRule implements IValidationRule<ImdgValidationContext
context.storeObject(ValidationStored.Sdf21Account, account.iterator().next());
return empty();
}
},
InstrumentPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<SDf21> context) {
SDf21 sdf21 = context.getValidatedObject();
if (TextUtil.isEmpty(sdf21.getSecurityCode())) {
return of(ClearingError.SecurityNotFound, sdf21.getSecurityCode());
}
SecuritySelector<Security> secSel = new SecuritySelector<>(
context.obtainMap(IMDGDistributedNames.Map_FixedIncomeSecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_MoneyMarketSecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_EquitySecurity, Security.class)
);
Security security = secSel.selectSecurityBySymbol(sdf21.getSecurityCode());
if (security == null) {
return of(ClearingError.SecurityNotFound, sdf21.getSecurityCode());
}
return empty();
}
}
;
};
@Override
public String ruleName() {

View file

@ -0,0 +1,47 @@
package ru.spcex.clearing.service.validation;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidationRule;
import java.util.Objects;
import java.util.Optional;
import java.util.function.Function;
public class SecurityBySecurityCodeValidationRule<T> implements IValidationRule<ImdgValidationContext<T>> {
private final Function<T, String> securitySymbolExtractor;
private SecurityBySecurityCodeValidationRule(Function<T, String> securitySymbolExtractor) {
Objects.requireNonNull(securitySymbolExtractor, "cannot create SecurityBySecurityCodeValidationRule: securitySymbolExtractor is null");
this.securitySymbolExtractor = securitySymbolExtractor;
}
public static <C> SecurityBySecurityCodeValidationRule<C> instance(Function<C, String> depoCodeExtractor) {
return new SecurityBySecurityCodeValidationRule<>(depoCodeExtractor);
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<T> context) {
T validatedObject = context.getValidatedObject();
String securitySymbol = securitySymbolExtractor.apply(validatedObject);
if (TextUtil.isEmpty(securitySymbol)) {
return of(ClearingError.SecurityNotFound, securitySymbol);
}
SecuritySelector<Security> secSel = new SecuritySelector<>(
context.obtainMap(IMDGDistributedNames.Map_FixedIncomeSecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_MoneyMarketSecurity, Security.class),
context.obtainMap(IMDGDistributedNames.Map_EquitySecurity, Security.class)
);
Security security = secSel.selectSecurityBySymbol(securitySymbol);
if (security == null) {
return of(ClearingError.SecurityNotFound, securitySymbol);
}
return empty();
}
}

View file

@ -10,9 +10,13 @@ public enum ValidationStored {
Sdf08Company, Sdf08Account,
Sdf21Company, Sdf21Account,
Sdf10Company, Sdf10Account,
Sdf21Account,
Sdf06Company, Sdf06Account,
ReturnDepositDmx
ReturnDepositDmx,
CompanyByDepoCode
}

View file

@ -10,6 +10,7 @@ public enum SdfTable implements IEnumKey {
SDF_08("SDF_08"),
SDF_09("SDF_09"),
SDF_10("SDF_10"),
SDF_11("SDF_11"),
SDF_12("SDF_12"),
SDF_13("SDF_13"),
SDF_16("SDF_16"),