diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java index 20cad7ef4..396ab021e 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java @@ -10,6 +10,10 @@ 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.ClearingAccount; +import ru.clearing.classes.statics.data.company.ClearingMemberCategory; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.company.relation.Relation; +import ru.clearing.classes.statics.data.sdf.SDf52; import ru.spcex.clearing.account.errors.AccountError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; @@ -26,22 +30,20 @@ import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.util.services.RequestHelper; import ru.spcex.clearing.validation.common.ValidationHelper; -import ru.spcex.platform.enumeration.AccountType; -import ru.spcex.platform.enumeration.SdfTable; -import ru.spcex.platform.enumeration.ServiceStatus; -import ru.spcex.platform.enumeration.WorkflowStatus; +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.ImdgTransaction; import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.utils.collection.Pair; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumKey; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.validation.IValidator; import java.time.Instant; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.function.Function; @Service @@ -56,6 +58,12 @@ public class ClearingAccountService extends QueueConsumer implements Initializin private final Function clearingAccountNewRequestValidator; private final Function clearingAccountUpdateRequestValidator; + private final Imdg sdf52Imdg; + private final Imdg accountImdg; + private final Imdg companyImdg; + private final Imdg relationImdg; + private final Imdg clearingMemberCategoryImdg; + @Autowired public ClearingAccountService(Consumer kafkaQueue, Producer kafkaResponseQueue, @@ -78,6 +86,12 @@ public class ClearingAccountService extends QueueConsumer implements Initializin this.requestHelper = requestHelper; this.clearingAccountNewRequestValidator = clearingAccountNewRequestValidator; this.clearingAccountUpdateRequestValidator = clearingAccountUpdateRequestValidator; + + this.sdf52Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf52, SDf52.class); + this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); + this.relationImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class); + this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class); } @Override @@ -91,6 +105,9 @@ public class ClearingAccountService extends QueueConsumer implements Initializin callback(AccountSdf01Request.class) .setFunction(this::accountNewSdf01) .forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put); + callback(StatementRequest.class) + .setFunction(this::accountUpdateSdf52) + .forDestination(Consts.ACCOUNT_PROCESS_SDF52, callbacks::put); init(); } @@ -106,7 +123,8 @@ public class ClearingAccountService extends QueueConsumer implements Initializin Instant now = Instant.now(); Account account = new Account(); account.setAccount(req.getAccount()); - account.setAccountType(AccountType.Clrn.getKey());if (req.getStatus() == null) { + account.setAccountType(AccountType.Clrn.getKey()); + if (req.getStatus() == null) { account.setStatus(WorkflowStatus.Active.getKey()); log.trace("Status not set in request. Use default: {}", account.getStatus()); } else { @@ -281,6 +299,130 @@ public class ClearingAccountService extends QueueConsumer implements Initializin return null; } + String serviceTypeByCMCOfCompany(Long companyId) { + Collection cmCategorys = clearingMemberCategoryImdg.getCollectionObjectsByFieldValues( + Map.of("companyId", companyId)); + String serviceType = null; + for (ClearingMemberCategory cmc : cmCategorys) { + if (IEnumKey.contains(cmc.getClearingMemberCategory(), ClearingCategory.B, ClearingCategory.I, ClearingCategory.V)) { + serviceType = ru.spcex.platform.enumeration.Service.MKR.getKey(); + } + if (IEnumKey.contains(cmc.getClearingMemberCategory(), ClearingCategory.F, ClearingCategory.C)) { + serviceType = ru.spcex.platform.enumeration.Service.MKR.getKey(); + } + } + return serviceType; + } + + /** + * @param status sdf.getStatus() + * @return AccountStatus или null + */ + AccountStatus parseSdf52Status(Long status) { + if (Long.valueOf(1L).equals(status)) { + return AccountStatus.ACTIVE; + } else if (Long.valueOf(0L).equals(status)) { + return AccountStatus.BLOCKED; + } + if (Long.valueOf(2L).equals(status)) { + return AccountStatus.CLOSE; + } + return null; + } + + public RequestInfoUpdate accountUpdateSdf52(BaseRequest systemRequest) { + log.debug("accountUpdateSdf52 StatementRequest received, id={}", systemRequest.getId()); + + StatementRequest req = systemRequest.getRequestPayload(); + Long groupId = req.getGroupId(); + if (SdfTable.SDF_52 != req.getTable()) { + log.warn("Unsupported table {} received on accountUpdateSdf52. Expected only {}.", + req.getTable(), SdfTable.SDF_52); + } + Collection sdfs = sdf52Imdg.getCollectionObjectsByFieldValues(Map.of( + "generationId", groupId + )); + if (sdfs.isEmpty()) { + log.info("Do not processing SDF52: S_DF52 not found by groupId={}", groupId); + return null; + } + log.info("start processing {} SDF52: groupId={}", sdfs.size(), groupId); + + List> toUpdate = new ArrayList<>(); + { // 1. Выборка данных + for (SDf52 sDf52 : sdfs) { + if (parseSdf52Status(sDf52.getStatus()) == null) { + String msg = messageResolver.resolve(new EnumMessage(AccountError.WrongFieldValue, "status")); + log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg); + continue; + } + Company company = sDf52.getDeal() == null ? null : companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", sDf52.getDeal())); + if (company == null) { + String msg = messageResolver.resolve(new EnumMessage(AccountError.CompanyNotFound, sDf52.getDeal())); + log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg); + continue; + } + String serviceType = serviceTypeByCMCOfCompany(company.getId()); + Relation relation = serviceType == null ? null : relationImdg.getSingleObjectByFieldValues(Map.of( + "consumerId", company.getId(), + "service", serviceType + )); + if (relation == null) { + String msg = messageResolver.resolve(new EnumMessage(AccountError.ClearingCategoryNotFound, company.getTradingCode(), serviceType)); + log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg); + continue; + } + Account account = sDf52.getAccount() == null ? null : accountImdg.getFirstObjectByFieldValues(Map.of( + "accountType", AccountType.Clrn.getKey(), + "account", sDf52.getAccount(), + "relationId", relation.getId() + )); + if (account == null) { + String msg = messageResolver.resolve(new EnumMessage(AccountError.AccountNotFound, sDf52.getAccount())); + log.warn("By generationId={} s_df52[{}] error: {}", groupId, sDf52.getId(), msg); + continue; + } + toUpdate.add(new Pair<>(sDf52, account)); + } + } + log.debug("Selected to update {} account's", toUpdate.size()); + int countOfUpdated = 0; + ImdgTransaction imdgTransaction = imdgProvider.newTransaction(); + boolean txOk = false; + imdgTransaction.beginTransaction(); + try { // 2. обновление данных, в транзакции + Imdg accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Instant now = Instant.now(); + accountsLoop: + for (Pair item : toUpdate) { + SDf52 sdf = item.getFirst(); + Account account = item.getSecond(); + AccountStatus newStatus = parseSdf52Status(sdf.getStatus()); + if (newStatus == null) { + throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus()); + } + if (!newStatus.equalsByKey(account.getStatus())) { + account.setStatus(newStatus.getKey()); + account.setUpdated(now); + accountImdg.update(account); + log.trace("S_DF52[{}] do update status to {} for account[{}]", sdf.getId(), newStatus.getKey(), account.getId()); + countOfUpdated++; + } + } + txOk = true; + } finally { + if (txOk) { + imdgTransaction.commitTransaction(); + } else { + log.debug("failed update, clearing accounts, rollback transaction. Request id={}; sdf52 groupId={}", systemRequest.getId(), groupId); + imdgTransaction.rollbackTransaction(); + } + } + log.debug("successfully processed, grouping id={}. Updated {} of {} accounts.", + groupId, countOfUpdated, toUpdate.size()); + return null; + } + public void sendStatementRequestBack(Long groupingSdf01Id, Long groupSdf02Id, List results) { StatementRequest request = new StatementRequest(); request.setGroupId(groupingSdf01Id); diff --git a/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/ClearingAccountServiceTest.java b/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/ClearingAccountServiceTest.java index 6c149f1dd..e416e16da 100644 --- a/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/ClearingAccountServiceTest.java +++ b/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/ClearingAccountServiceTest.java @@ -4,6 +4,7 @@ import org.apache.kafka.clients.consumer.MockConsumer; import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; +import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.ArgumentCaptor; @@ -224,6 +225,15 @@ class ClearingAccountServiceTest { ACCOUNT_MATCHER.assertMatch(accountResult, predictableAccount); } + @Test + void parseSdf52Status() { + Assertions.assertEquals(AccountStatus.BLOCKED, clearingAccountService.parseSdf52Status(0L)); + Assertions.assertEquals(AccountStatus.ACTIVE, clearingAccountService.parseSdf52Status(1L)); + Assertions.assertEquals(AccountStatus.CLOSE, clearingAccountService.parseSdf52Status(2L)); + Assertions.assertEquals(null, clearingAccountService.parseSdf52Status(3L)); + Assertions.assertEquals(null, clearingAccountService.parseSdf52Status(null)); + } + // /** // * {@link ClearingAccountService#accountNewSdf01(BaseRequest)}
// * Тест проверяет создание сущности {@link BaseRequest} в Hazelcast при передаче из Apache Kafka.
diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java index 742e95a69..c5f42d148 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java @@ -34,6 +34,7 @@ public class DbfImportKafkaMessenger implements InitializingBean { messengers.put(ETable.DF_06, groupId -> messageStatement(groupId, SdfTable.SDF_06, Consts.STATEMENT_PROCESS_SDF06)); messengers.put(ETable.DF_09, groupId -> messageStatement(groupId, SdfTable.SDF_09)); messengers.put(ETable.DF_16, groupId -> messageStatement(groupId, SdfTable.SDF_16)); + messengers.put(ETable.DF_52, groupId -> messageStatement(groupId, SdfTable.SDF_52, Consts.ACCOUNT_PROCESS_SDF52)); messengers.put(ETable.DF_55, groupId -> messageStatement(groupId, SdfTable.SDF_55)); messengers.put(ETable.DF_57, groupId -> messageStatement(groupId, SdfTable.SDF_57)); } diff --git a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java index 8c228aef8..21e2847ec 100644 --- a/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java +++ b/platform-parent/platform-enum/src/main/java/ru/spcex/platform/enumeration/SdfTable.java @@ -16,6 +16,7 @@ public enum SdfTable implements IEnumKey { SDF_16("SDF_16"), SDF_20("SDF_20"), SDF_21("SDF_21"), + SDF_52("SDF_52"), SDF_55("SDF_55"), SDF_57("SDF_57"); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 0e2b55d3a..0ae14d291 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -90,6 +90,7 @@ public interface Consts { String ACCOUNT_NEW_SDF01 = "account-new-sdf01"; String ACCOUNT_NEW_SDF08 = "account-new-sdf08"; + String ACCOUNT_PROCESS_SDF52 = "account-process-sdf52"; String DESTINATION_RELATION_NEW = "relation-new"; String DESTINATION_RELATION_UPDATE = "relation-update";