account-service TCR validation bugfix

clearing-service поменял местами SDF01Executor и SDF57Executor
разнес нотификацию сверки и Sdf01Executor
This commit is contained in:
ialbert 2023-07-06 14:52:27 +03:00
parent 30349030dd
commit 6c1ec5380e
8 changed files with 197 additions and 37 deletions

View file

@ -1,6 +1,5 @@
package ru.spcex.clearing.account.config.validation;
import com.hazelcast.query.PredicateBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.clearing.classes.statics.data.account.Account;
@ -31,7 +30,6 @@ import ru.spcex.platform.utils.validation.IValidator;
import ru.spcex.platform.utils.validation.ValidatorImpl;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.function.Consumer;
@ -144,7 +142,7 @@ public class TradingClearingRegistryValidationConfig {
TradingClearingRegistryNewRequest validatedObject = context.getValidatedObject();
Imdg<TradingClearingRegistry> tcrMap = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
if (validatedObject.getMoneyAccountId() != null) {
if (validatedObject.getMoneyAccountId() == null) {
return of(AccountError.RequiredFieldEmpty, "MoneyAccountId");
}
//Map<String, Comparable<?>> query = new HashMap<>();

View file

@ -273,6 +273,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
StatementRequest request = new StatementRequest();
request.setGroupId(groupingSdf01Id);
request.setAccountCreationResults(results);
request.setContinueSdf01(true);
request.setTable(SdfTable.SDF_01); // по нему запрос получили
log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request));
kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request);

View file

@ -211,6 +211,7 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea
public void sendStatementRequestBack(Long groupingSdf01Id, List<AccountSdfToStatementRequestPart> results) {
StatementRequest request = new StatementRequest();
request.setGroupId(groupingSdf01Id);
request.setContinueSdf01(true); //fixme????
request.setAccountCreationResults(results);
request.setTable(SdfTable.SDF_08); // по нему запрос получили
log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request));

View file

@ -108,7 +108,7 @@ public class RegistryService {
prdBldr.or(moneyPredicate.orElse(prdBldr.alwaysTrue()), depoPredicate.orElse(prdBldr.alwaysTrue())),
prdBldr.equals("companyId", companyId),
prdBldr.equals("accountId", moneyAccountId),
prdBldr.sql("tradingClearingRegistryId is null")
prdBldr.sql("tradingClearingRegistryId == null")
));
}

View file

@ -18,6 +18,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.ContinueSessionB
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.executors.AbstractExecutor;
import ru.spcex.clearing.service.executors.Reviser;
import ru.spcex.clearing.service.model.Result;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.SdfTable;
@ -36,6 +37,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
private final KafkaSender kafkaSender;
private final Map<SdfTable, Imdg<? extends SpcexObjectBase>> sdfImdgs;
private final Map<SdfTable, AbstractExecutor<?>> executorsMap;
private final Reviser reviser;
/**
* Пары sdf запросов пришедшие с модуля dbf-import
@ -46,9 +48,10 @@ public class StatementService extends QueueConsumer implements InitializingBean
public StatementService(Consumer<String, Object> kafkaQueue,
ImdgProvider imdgProvider,
KafkaSender kafkaSender,
@Qualifier("sdfExecutors") Map<SdfTable, AbstractExecutor<?>> executorsMap) {
@Qualifier("sdfExecutors") Map<SdfTable, AbstractExecutor<?>> executorsMap, Reviser reviser) {
super(kafkaQueue);
this.imdgProvider = imdgProvider;
this.reviser = reviser;
this.sdfImdgs = new EnumMap<>(SdfTable.class);
this.kafkaSender = kafkaSender;
this.executorsMap = executorsMap;
@ -62,11 +65,47 @@ public class StatementService extends QueueConsumer implements InitializingBean
@Override
public void afterPropertiesSet() throws Exception {
callback(StatementRequest.class)
.setConsumer(this::processPaired)
.setConsumer(systemRequest -> {
if (systemRequest.getRequestPayload() != null && systemRequest.getRequestPayload().isContinueSdf01()) {
Optional<Long> sdf01And57Key = findCompletePair();
if (sdf01And57Key.isEmpty()) {
log.error("FATAL: response from account-service received id={} sdf ids={}", systemRequest.getId(),
systemRequest.getRequestPayload()
.getAccountCreationResults()
.stream()
.map(res -> String.valueOf(res.getSdfId()))
.collect(Collectors.joining(",", "[", "]"))
);
return;
}
processSdf01And57(sdf01And57Key.get(), systemRequest.getRequestPayload());
log.info("Sdf01 and Sdf57 processed successfully after account-service command.");
} else {
processPaired(systemRequest);
}
})
.forDestination(Consts.STATEMENT_PROCESS, callbacks::put);
init();
}
private void processSdf01And57(Long key, StatementRequest fromAccService) {
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
boolean allProcessed = processSdf01(fromAccService == null ? pair.getFirst() : fromAccService);
if (!allProcessed) {
log.info("sdf01 execution wasn't complete, waiting for an answer from account-service");
return;
}
//затем sdf57
processSdf57(pair.getSecond());
reviser.doRevise(pair.getFirst().getGroupId());
//теперь можем продолжить сессию с шага 1
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
pairOfSdfRequest.remove(key);
log.info("pair sdf01/sdf57 processed successfully");
}
private void processPaired(BaseRequest<StatementRequest> systemRequest) {
StatementRequest statementRequest = systemRequest.getRequestPayload();
SdfTable table = statementRequest.getTable();
@ -77,18 +116,9 @@ public class StatementService extends QueueConsumer implements InitializingBean
boolean doSomeone = false;
if (List.of(SdfTable.SDF_01, SdfTable.SDF_57).contains(table)) {
{
//всегда сначала обработаем sdf57
Long key = completePairKey.get();
Pair<StatementRequest, StatementRequest> pair = pairOfSdfRequest.get(key);
processSdf57(pair.getSecond());
//затем sdf01
processSdf01(pair.getFirst());
//теперь можем продолжить сессию с шага 1
ContinueSessionBnRequest continueSessionBn = new ContinueSessionBnRequest();
kafkaSender.sendRequestToQueue(Consts.CONTINUE_SESSION_BN_FIRST_PART, continueSessionBn);
pairOfSdfRequest.remove(key);
doSomeone = true;
log.info("sdf01/sdf57 pair is received");
processSdf01And57(completePairKey.get(), null);
return;
}
} else if (List.of(SdfTable.SDF_08, SdfTable.SDF_04).contains(table)) {
if (table == SdfTable.SDF_08) {
@ -166,7 +196,7 @@ public class StatementService extends QueueConsumer implements InitializingBean
// finishSendCommand(res, service, statementRequest);
}
private void processSdf01(StatementRequest statementRequest) {
private boolean processSdf01(StatementRequest statementRequest) {
Imdg<SDf01> sdfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
Collection<? extends SpcexObjectBase> sdfGroup;
if (statementRequest.getAccountCreationResults().size() == 0) {
@ -181,11 +211,15 @@ public class StatementService extends QueueConsumer implements InitializingBean
.collect(Collectors.toList());
}
AbstractExecutor service = executorsMap.get(SdfTable.SDF_01);
Result res = service.execute(sdfGroup, statementRequest); if (res.getAccountRequests().size() != 0) {
Result res = service.execute(sdfGroup, statementRequest);
if (res.getAccountRequests().size() != 0) {
kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests()));
return false;
} else if (service.isNeedToSendCommand()) {
service.sendCommand(kafkaSender, res);
return true;
}
return true;
}
private Optional<Long> saveRequest(StatementRequest statementRequest) {
@ -248,6 +282,16 @@ public class StatementService extends QueueConsumer implements InitializingBean
return uncompletedPair;
}
private Optional<Long> findCompletePair() {
return pairOfSdfRequest.entrySet()
.stream()
.filter(entry -> entry.getValue().getFirst() != null && entry.getValue().getSecond() != null)
.filter(entry -> SdfTable.SDF_01.equals(entry.getValue().getFirst().getTable())
&& SdfTable.SDF_57.equals(entry.getValue().getSecond().getTable()))
.map(Map.Entry::getKey)
.findFirst();
}
private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdfRequestPart> accountRequests) {
AccountSdf01Request r = new AccountSdf01Request();
r.setGroupingSdf01Id(sdf01GroupingId);

View file

@ -0,0 +1,107 @@
package ru.spcex.clearing.service.executors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.sdf.SDf01;
import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
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.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.number.BigDecimalUtil;
import java.math.BigDecimal;
import java.util.Collection;
import java.util.Optional;
@Component
public class Reviser {
private final Logger log = LoggerFactory.getLogger(Reviser.class);
private static final String reviseFailedMessage = "Сверка остатков денежных средств по результатам клиринговой сессии завершена с ошибками.";
private static final String reviseSuccessMessage = "Сверка остатков денежных средств по результатам клиринговой сессии завершена успешно.";
private final Imdg<Registry> registryImdg;
private final Imdg<Statement> statementImdg;
private final Imdg<SDf01> sdf01Imdg;
private final KafkaSender kafkaSender;
public Reviser(ImdgProvider imdgProvider, KafkaSender kafkaSender) {
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.sdf01Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
this.kafkaSender = kafkaSender;
}
public void doRevise(Long sdf01Group) {
Collection<SDf01> sdfs = findSdf01ByGroupId(sdf01Group);
log.info("doing revise for {} sdfs", sdfs.size());
boolean reviseFailed = false;
for (SDf01 sdf : sdfs) {
Optional<Statement> statement = statementsBySdf01(sdf);
if (statement.isEmpty()) {
log.trace("no statement for sdf01.id={}", sdf.getId());
continue;
}
Optional<Registry> registry = findReg(statement.get());
if (registry.isEmpty()) {
log.trace("no registry for sdf01.id={} statement.id={}", sdf.getId(), statement.get().getId());
continue;
}
if (!checkDiffBalance(registry.get().getDiffBalance())) {
log.info("sdf01.id={}, stmt.id={}, registry.id={} diffBalance is not zero: {}",
sdf.getId(),
statement.get().getId(),
registry.get().getId(),
registry.get().getDiffBalance());
reviseFailed = true;
}
log.trace("revise ok for sdf01.id={}, stmt.id={}, registry.id={}", sdf.getId(), statement.get().getId(), registry.get().getId());
}
NotificationNewRequest reviseNotification = new NotificationNewRequest();
reviseNotification.setObjectType(ObjectType.rgst.getKey());
reviseNotification.setComment(reviseFailed ? reviseFailedMessage : reviseSuccessMessage);
reviseNotification.setPriority(reviseFailed ? Priority.HIGH.getKey() : Priority.LOW.getKey());
kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, reviseNotification);
log.info("revise notification sent (revise {})", reviseFailed ? "error" : "success");
}
private boolean checkDiffBalance(BigDecimal diffBalance) {
return BigDecimalUtil.safeBD(diffBalance).compareTo(BigDecimal.ZERO) == 0;
}
private Collection<SDf01> findSdf01ByGroupId(Long groupId) {
return sdf01Imdg.getCollectionObjectsBySQL("generationId = " + groupId);
}
private Optional<Statement> statementsBySdf01(SDf01 sdf01) {
ImdgPredicateBuilder pb = statementImdg.predicateBuilder();
ImdgPredicate stmtPredicate = pb.and(
pb.equals("inSDfId", sdf01.getId()),
pb.equals("inOutSDfType", InOutSDfType.type1.getKey())
);
return Optional.ofNullable(statementImdg.getSingleObjectByPredicate(stmtPredicate));
}
private Optional<Registry> findReg(Statement s) {
RegistryTradingParams p = new RegistryTradingParams(
RegistryDesignation.A, RegistryInstrumentType.M, null, RegistryUnit.T
);
String sql = RegistryCodeSqlBuilder.getInstance(p).build();
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql),
pb.sql(sql),
pb.equals("accountId", s.getAccountId()),
pb.equals("companyId", s.getAddresseeId())
);
return Optional.ofNullable(registryImdg.getSingleObjectByPredicate(rgstrPredicate));
}
}

View file

@ -21,7 +21,6 @@ import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.LoggingService;
import ru.spcex.clearing.service.model.Result;
@ -106,7 +105,7 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
Result result = new Result();
Long generationIdForGroup = imdgProvider.getImdgIdGenerator().nextId();
result.setGenerationId(generationIdForGroup);
boolean reviseFailed = false;
// boolean reviseFailed = false;
log.info("SDF01 execution: sdf01 number={}, groupId={}", sdf.size(), sdf.stream().findFirst().map(SDf01::getGenerationId).orElse(null));
for (SDf01 sdf01 : sdf) {
IValidator validator = sDf01Validator.apply(sdf01);
@ -160,14 +159,14 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
rgs = createRegistryByStatement(stmt, company, account);
registryImdg.insert(rgs);
}
if (!checkDiffBalance(rgs.getDiffBalance())) {
log.info("sdf01.id={}, stmt.id={}, registry.id={} diffBalance is not zero: {}",
sdf01.getId(),
stmt.getId(),
rgs.getId(),
rgs.getDiffBalance());
reviseFailed = true;
}
//if (!checkDiffBalance(rgs.getDiffBalance())) {
// log.info("sdf01.id={}, stmt.id={}, registry.id={} diffBalance is not zero: {}",
// sdf01.getId(),
// stmt.getId(),
// rgs.getId(),
// rgs.getDiffBalance());
// reviseFailed = true;
//}
stmt.setOperationStatus(OperationStatus.Executed.getKey());
statementImdg.update(stmt);
} else {
@ -177,12 +176,12 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
statementImdg.update(stmt);
}
}
NotificationNewRequest reviseNotification = new NotificationNewRequest();
reviseNotification.setObjectType(ObjectType.rgst.getKey());
reviseNotification.setComment(reviseFailed ? reviseFailedMessage : reviseSuccessMessage);
reviseNotification.setPriority(reviseFailed ? Priority.HIGH.getKey() : Priority.LOW.getKey());
kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, reviseNotification);
log.info("revise notification sent (revise {})", reviseFailed ? "error" : "success");
//NotificationNewRequest reviseNotification = new NotificationNewRequest();
//reviseNotification.setObjectType(ObjectType.rgst.getKey());
//reviseNotification.setComment(reviseFailed ? reviseFailedMessage : reviseSuccessMessage);
//reviseNotification.setPriority(reviseFailed ? Priority.HIGH.getKey() : Priority.LOW.getKey());
//kafkaSender.sendRequestToQueue(Consts.NOTIFICATION_NEW, reviseNotification);
//log.info("revise notification sent (revise {})", reviseFailed ? "error" : "success");
return result;
}
@ -291,7 +290,7 @@ public class Sdf01Executor extends AbstractExecutor<SDf01> {
rgs.setAccount(account.getAccount());
rgs.setRegistryDesignation(RegistryDesignation.A.getKey());
rgs.setRegistryInstrumentType(RegistryInstrumentType.M.getKey());
ClearingAccount accountForStatement = clearingAccountImdg.getSingleObjectByID(statement.getAccountId());
ClearingAccount accountForStatement = clearingAccountImdg.getSingleObjectByFieldValues(Map.of("accountId", statement.getAccountId()));
if (accountForStatement != null) {
rgs.setRegistryCapacity(accountForStatement.getClearingAccountType());
}

View file

@ -15,11 +15,21 @@ public class StatementRequest {
//from account-service creation
@JsonProperty
List<AccountSdfToStatementRequestPart> accountCreationResults = new ArrayList<>();
@JsonProperty
private boolean continueSdf01 = false;
public Long getGroupId() {
return groupId;
}
public boolean isContinueSdf01() {
return continueSdf01;
}
public void setContinueSdf01(boolean continueSdf01) {
this.continueSdf01 = continueSdf01;
}
public void setGroupId(Long groupId) {
this.groupId = groupId;
}