diff --git a/clearing-parent/account-service/pom.xml b/clearing-parent/account-service/pom.xml index ddbf3d9d1..aa06e621d 100644 --- a/clearing-parent/account-service/pom.xml +++ b/clearing-parent/account-service/pom.xml @@ -24,6 +24,10 @@ ru.spcex.platform platform-imdg-api-hazelcast-impl + + ru.spcex.platform + platform-enum + ru.spcex.clearing classes diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/ErrorResolverConfig.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/ErrorResolverConfig.java new file mode 100644 index 000000000..7cf2db17a --- /dev/null +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/ErrorResolverConfig.java @@ -0,0 +1,14 @@ +package ru.spcex.clearing.account.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.enumeration.SimpleMessageResolver; + +@Configuration +public class ErrorResolverConfig { + @Bean + public IMessageResolver messageResolver() { + return new SimpleMessageResolver(); + } +} diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/errors/AccountError.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/errors/AccountError.java new file mode 100644 index 000000000..76dc0ac4b --- /dev/null +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/errors/AccountError.java @@ -0,0 +1,17 @@ +package ru.spcex.clearing.account.errors; + +import ru.spcex.platform.utils.enumeration.IEnumId; + +public enum AccountError implements IEnumId { + AccountAlreadyExist(5010L); + private final Long id; + + AccountError(Long id) { + this.id = id; + } + + @Override + public Long getId() { + return id; + } +} diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java index 3c1ad9ed6..cd8817ca8 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java @@ -7,7 +7,10 @@ import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.BankAccount; +import ru.clearing.classes.statics.data.company.relation.Relation; +import ru.spcex.clearing.account.errors.AccountError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; @@ -15,36 +18,67 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewReq import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; +import ru.spcex.platform.enumeration.AccountStatus; +import ru.spcex.platform.enumeration.AccountType; +import ru.spcex.platform.enumeration.Allowed; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IMessageResolver; + +import java.time.Instant; +import java.util.Collection; @Service public class BankAccountService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg bankAccountMap; + private final Imdg accountMap; + private final Imdg relationMap; + private final IMessageResolver messageResolver; @Autowired - public BankAccountService(Consumer kafkaQueue, Producer kafkaProducer, ImdgProvider imdgProvider) { + public BankAccountService(Consumer kafkaQueue, + Producer kafkaProducer, + ImdgProvider imdgProvider, + IMessageResolver messageResolver) { super(kafkaQueue, kafkaProducer); this.bankAccountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_BankAccount, BankAccount.class); + this.accountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + this.relationMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class); + this.messageResolver = messageResolver; } @Override public void afterPropertiesSet() { callback(BankAccountNewRequest.class) - .setConsumer(this::bankAccountNew) + .setFunction(this::bankAccountNew) .forDestination(Consts.DESTINATION_BANK_ACCOUNT_NEW, callbacks::put); callback(BankAccountUpdateRequest.class) - .setConsumer(this::bankAccountUpdate) + .setFunction(this::bankAccountUpdate) .forDestination(Consts.DESTINATION_BANK_ACCOUNT_UPDATE, callbacks::put); callback(CommonDeleteRequest.class) - .setConsumer(this::bankAccountDelete) + .setFunction(this::bankAccountDelete) .forDestination(Consts.DESTINATION_BANK_ACCOUNT_DELETE, callbacks::put); init(); } - private void bankAccountNew(BaseRequest userRequest) { + private RequestInfoUpdate bankAccountNew(BaseRequest userRequest) { BankAccountNewRequest req = userRequest.getRequestPayload(); + + Collection accountsByKey = + accountMap.getCollectionObjectsBySQL(String.format("account = %s", req.account)); + if (!accountsByKey.isEmpty()) { + String errorMsg = messageResolver.resolve(new EnumMessage(AccountError.AccountAlreadyExist)); + log.error("cannot process MoneyMarketSecurityNewRequest id={}: {}", userRequest.getId(), errorMsg); + return new RequestInfoUpdate() + .setId(userRequest.getId()) + .setStatus(Status.Error) + .setMessage(errorMsg); + } + log.debug("BankAccountNewRequest received"); BankAccount bankAccount = new BankAccount(); bankAccount.setBankIdentificationCode(req.getBankIdentificationCode()); @@ -57,11 +91,33 @@ public class BankAccountService extends QueueConsumer implements InitializingBea bankAccount.setTaxRegistrationReasonCode(req.getTaxRegistrationReasonCode()); bankAccount.setAccount(req.getAccount()); + Account account = new Account(); + account.setAccount(req.account); + account.setAccountType(AccountType.Bank.getKey()); + + String relationSqlCondition = String.format("consumerId = %s and service = %s", req.companyId, + ru.spcex.platform.enumeration.Service.MKR.getKey()); + Relation relationByCompany = relationMap.getSingleObjectBySQL(relationSqlCondition); + if (relationByCompany != null) { + account.setRelationId(relationByCompany.getId()); + account.setCompanyId(relationByCompany.getConsumerId()); + } else { + log.warn("Not found relation by condition: {}", relationSqlCondition); + } + account.setAccountStatus(AccountStatus.ACTIVE.getKey()); + account.setProcessingSign(Allowed.ALLOWED.getKey()); + account.setCreated(Instant.now()); + account.setUpdated(Instant.now()); + + accountMap.insert(account); + + bankAccount.setAccountId(account.getId()); bankAccountMap.insert(bankAccount); log.debug("successfully processed, new id {}", bankAccount.getId()); + return null; } - private void bankAccountUpdate(BaseRequest userRequest) { + private RequestInfoUpdate bankAccountUpdate(BaseRequest userRequest) { BankAccountUpdateRequest req = userRequest.getRequestPayload(); log.debug("BankAccountUpdateRequest received id = {}", req.getId()); BankAccount bankAccount = bankAccountMap.getSingleObjectByID(req.getId()); @@ -75,17 +131,29 @@ public class BankAccountService extends QueueConsumer implements InitializingBea bankAccount.setTaxRegistrationReasonCode(req.getTaxRegistrationReasonCode()); bankAccount.setAccount(req.getAccount()); + Account account = accountMap.getSingleObjectByID(bankAccount.getAccountId()); + account.setAccount(req.account); + account.setUpdated(Instant.now()); + + accountMap.update(account); bankAccountMap.update(bankAccount); log.debug("successfully update, existing bankAccount with id {}", bankAccount.getId()); + return null; } - private void bankAccountDelete(BaseRequest userRequest) { + private RequestInfoUpdate bankAccountDelete(BaseRequest userRequest) { CommonDeleteRequest req = userRequest.getRequestPayload(); log.debug("CommonDeleteRequest received id = {}", req.getId()); BankAccount bankAccount = bankAccountMap.getSingleObjectByID(req.getId()); + Account account = accountMap.getSingleObjectByID(bankAccount.getAccountId()); + account.setAccountStatus(AccountStatus.BLOCKED.getKey()); + account.setUpdated(Instant.now()); + + accountMap.update(account); bankAccountMap.delete(bankAccount); log.debug("successfully delete, existing bankAccount with id {}", bankAccount.getId()); + return null; } } diff --git a/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/BankAccountServiceTest.java b/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/BankAccountServiceTest.java index 6060126ce..ed1a44667 100644 --- a/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/BankAccountServiceTest.java +++ b/clearing-parent/account-service/src/test/java/ru/spcex/clearing/account/service/BankAccountServiceTest.java @@ -17,6 +17,7 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit.jupiter.SpringExtension; import ru.clearing.classes.statics.data.account.BankAccount; +import ru.spcex.clearing.account.config.ErrorResolverConfig; import ru.spcex.clearing.account.config.HazelcastServiceTestConfiguration; import ru.spcex.clearing.account.utils.MatcherFactory.Matcher; import ru.spcex.clearing.imdg.IMDGDistributedNames; @@ -27,6 +28,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountNewReq import ru.spcex.clearing.platform.messaging.domain.cud.account.BankAccountUpdateRequest; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest; import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService; +import ru.spcex.platform.utils.enumeration.IMessageResolver; import java.util.Collections; import java.util.HashMap; @@ -35,7 +37,7 @@ import static ru.spcex.clearing.account.utils.MatcherFactory.usingIgnoringFields @ExtendWith(SpringExtension.class) @ContextConfiguration(classes = { - HazelcastServiceTestConfiguration.class}) + HazelcastServiceTestConfiguration.class, ErrorResolverConfig.class}) public class BankAccountServiceTest { public static final Matcher BANK_ACCOUNT_MATCHER = usingIgnoringFieldsComparator(); private static final int PARTITION = 0; @@ -47,6 +49,8 @@ public class BankAccountServiceTest { @Autowired @Qualifier("hazelcastServiceTest") private HazelcastService hazelcastServiceTest; + @Autowired + private IMessageResolver messageResolver; private MockConsumer mockConsumer; private MockProducer mockProducer; @@ -120,7 +124,8 @@ public class BankAccountServiceTest { //ACT //service set up - BankAccountService bankAccountService = new BankAccountService(mockConsumer, mockProducer, hazelcastServiceTest); + BankAccountService bankAccountService = new BankAccountService(mockConsumer, mockProducer, + hazelcastServiceTest, messageResolver); Thread.sleep(10000); //callbacks set up bankAccountService.afterPropertiesSet(); @@ -221,7 +226,8 @@ public class BankAccountServiceTest { //ACT //service set up - BankAccountService bankAccountService = new BankAccountService(mockConsumer, mockProducer, hazelcastServiceTest); + BankAccountService bankAccountService = new BankAccountService(mockConsumer, mockProducer, + hazelcastServiceTest, messageResolver); Thread.sleep(10000); //callbacks set up bankAccountService.afterPropertiesSet(); @@ -305,7 +311,8 @@ public class BankAccountServiceTest { //ACT //service set up - BankAccountService bankAccountService = new BankAccountService(mockConsumer, mockProducer, hazelcastServiceTest); + BankAccountService bankAccountService = new BankAccountService(mockConsumer, mockProducer, + hazelcastServiceTest, messageResolver); Thread.sleep(10000); //callbacks set up bankAccountService.afterPropertiesSet(); diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java index f5429d6c9..a6e7bb993 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/scheduler/LauncherController.java @@ -81,7 +81,7 @@ public class LauncherController extends AbstractQueueController { } launcherCommand.setTask(dictionaryName); launcherCommand.setUserId(user.getId()); - //пока здесь + // пока здесь, это требуется для сохранения истории saveLauncher(dictionaryName, user.getId()); //топики ограничиваются наличием в taskDictionary //подписываются на разные топики в разных модулях, см. ru.spcex.platform.enumeration.Task#topic diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/schedule/LauncherNew.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/schedule/LauncherNew.java index 76b301904..ce810d0e8 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/schedule/LauncherNew.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/schedule/LauncherNew.java @@ -6,14 +6,16 @@ import io.swagger.annotations.ApiModelProperty; import ru.spcex.clearing.backendapi.domain.actions.IAction; import ru.spcex.clearing.backendapi.errors.BackEndError; import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportWithPeriodRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.utils.enumeration.EnumMessage; import java.util.Collection; import java.util.Collections; import java.util.List; -public class LauncherNew implements IAction { +public class LauncherNew implements IAction { @ApiModelProperty(value = "Идентификатор единоличного исполнительного органа", example = "ABCD") @JsonProperty private String task; @@ -22,11 +24,17 @@ public class LauncherNew implements IAction { private Long userId; @Override - public LauncherCommandRequest toRequest() { - LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest(); - taskRunnerCommandRequest.setTaskName(task); - taskRunnerCommandRequest.setUserId(userId); - return taskRunnerCommandRequest; + public Object toRequest() { + if (Task.createReport_GREP.equalsByKey(task)) { + ReportWithPeriodRequest reportCommand = new ReportWithPeriodRequest(); + reportCommand.setReportId(task); + return reportCommand; + } else { + LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest(); + taskRunnerCommandRequest.setTaskName(task); + taskRunnerCommandRequest.setUserId(userId); + return taskRunnerCommandRequest; + } } @Override diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/errors/BalanceError.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/errors/BalanceError.java index bc93391a6..773623c19 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/errors/BalanceError.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/errors/BalanceError.java @@ -7,7 +7,7 @@ public enum BalanceError implements IEnumId { CurrencyNotFound(5213L), CurrentDateOnly(5214L), WrongMarket(5215L), - WrongAccount(5215L), + WrongAccount(5216L), AccountNotPresent(5217L), AccountNotActive(5218L), BalanceNotEnough(5222L), diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf09Executor.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf09Executor.java index 96e8608bc..2291363ad 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf09Executor.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/service/Sdf09Executor.java @@ -137,6 +137,9 @@ public class Sdf09Executor extends AbstractExecutor { statement.setInOutDirection(InOutDirection.in.getKey()); statement.setAmount(sdf09.getSum()); statement.setCashMovementCurrencyCode(CurrencyCode.RUB.getKey()); + statement.setInSDfId(sdf09.getId()); +// statement.setOutSDfId(sdf10.getId()); заполняется внешним кодом + statement.setInOutSDfType(InOutSDfType.type9.getKey()); statementImdg.update(statement); } diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf01ValidationRule.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf01ValidationRule.java index c783e57d3..120cb137a 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf01ValidationRule.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf01ValidationRule.java @@ -58,7 +58,6 @@ public enum Sdf01ValidationRule implements IValidationRule validate(ImdgValidationContext context) { SDf01 sdf01 = context.getValidatedObject(); - //fixme string format??? if (!LocalDate.now().equals(LocalDate.parse(sdf01.getDat(), datFormatter))) { return of(BalanceError.CurrentDateOnly); } diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf09ValidationRule.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf09ValidationRule.java index 5daf7d4a6..325e00892 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf09ValidationRule.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf09ValidationRule.java @@ -1,5 +1,6 @@ package ru.spcex.clearing.balance.validation; +import org.slf4j.LoggerFactory; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.CompanySymbols; @@ -13,6 +14,7 @@ import ru.spcex.platform.imdg.validation.ImdgValidationContext; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.validation.IValidationRule; +import java.math.RoundingMode; import java.util.Map; import java.util.Optional; @@ -23,11 +25,17 @@ public enum Sdf09ValidationRule implements IValidationRule companySymbolImdg = context.obtainMap(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); Imdg companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class); + if (sdf09.getInn() == null) { + return of( BalanceError.CompanyNotFound, "empty inn"); + } + String innStr = sdf09.getInn().setScale(0, RoundingMode.DOWN).toString(); CompanySymbols symbol = companySymbolImdg .getSingleObjectByFieldValues(Map.of("companySymbol", CompanySymbol.INN.getKey(), - "companySymbolValue", sdf09.getInn())); - Company company = companyImdg.getSingleObjectByID(symbol.getCompanyId()); - + "companySymbolValue", innStr)); + Company company = null; + if (symbol != null) { + company = companyImdg.getSingleObjectByID(symbol.getCompanyId()); + } if (company == null) { return of( BalanceError.CompanyNotFound); } diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRule.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRule.java index 69bdeb806..fac3d544b 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRule.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRule.java @@ -1,5 +1,7 @@ package ru.spcex.clearing.balance.validation; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.CompanySymbols; import ru.clearing.classes.statics.data.sdf.SDf16; @@ -11,6 +13,7 @@ import ru.spcex.platform.imdg.validation.ImdgValidationContext; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.validation.IValidationRule; +import java.math.RoundingMode; import java.util.Optional; public enum Sdf16ValidationRule implements IValidationRule> { @@ -20,11 +23,26 @@ public enum Sdf16ValidationRule implements IValidationRule companySymbolsImdg = context.obtainMap(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); Imdg companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class); - CompanySymbols companySymbols = companySymbolsImdg.getSingleObjectBySQL( - String.format("(companySymbolValue = '%s' and companySymbol = '%s') " + - "or (companySymbolValue = '%s' and companySymbol = '%s')", - sdf16.getInn(), CompanySymbol.INN.getKey(), sdf16.getBic(), CompanySymbol.BIC.getKey())); + String sql = ""; + if (sdf16.getInn() != null) { + String innStr = sdf16.getInn().setScale(0, RoundingMode.DOWN).toString(); + sql = String.format("(companySymbolValue = '%s' and companySymbol = '%s') ", + innStr, CompanySymbol.INN.getKey()); + } + if (sdf16.getBic() != null) { + if (!sql.isEmpty()) sql += " or "; + String bicStr = sdf16.getBic().setScale(0, RoundingMode.DOWN).toString(); + sql += String.format("(companySymbolValue = '%s' and companySymbol = '%s') ", + bicStr, CompanySymbol.BIC.getKey()); + } + Logger log = LoggerFactory.getLogger(getClass()); + log.debug("SQL for search company by sdf16: {}", sql); + if (sql.isEmpty()) { + return of(BalanceError.CompanyNotFound, "inn and bic is empty"); + } Company company = null; + CompanySymbols companySymbols = companySymbolsImdg.getSingleObjectBySQL(sql); + log.debug("By sdf16 found companySymbol: {}", companySymbols == null ? null : companySymbols.getId()); if (companySymbols != null) { company = companyImdg.getSingleObjectByID(companySymbols.getCompanyId()); } diff --git a/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRuleTest.java b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRuleTest.java new file mode 100644 index 000000000..bd833cfc1 --- /dev/null +++ b/clearing-parent/balance-service/src/test/java/ru/spcex/clearing/balance/validation/Sdf16ValidationRuleTest.java @@ -0,0 +1,44 @@ +package ru.spcex.clearing.balance.validation; + +import org.junit.jupiter.api.Test; +import ru.clearing.classes.statics.data.sdf.SDf16; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.platform.classes.base.SpcexObjectBase; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.utils.enumeration.EnumMessage; + +import java.math.BigDecimal; +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.*; + +class Sdf16ValidationRuleTest { + + @Test + void values() { + ImdgValidationContext ctx = new ImdgValidationContext<>(); + SDf16 object = new SDf16(); +// object.setBic(BigDecimal.valueOf(1000004L)); + object.setInn(BigDecimal.valueOf(1000024L)); + ctx.setValidatedObject(object); + ctx.addImdg(IMDGDistributedNames.Map_CompanySymbols, makeTestMap()); + ctx.addImdg(IMDGDistributedNames.Map_Company, makeTestMap()); + Optional err = Sdf16ValidationRule.CompanyPresent.validate(ctx); + assertTrue(err.isPresent()); + } + + Imdg makeTestMap() { + return new Imdg() { + @Override + public String getMapName() { + return Imdg.super.getMapName(); + } + + @Override + public T getSingleObjectBySQL(String sql) { + return null; + } + }; + } +} \ No newline at end of file diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java index c67363f54..a2bebeafc 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/logic/stages/ExportFromHazelcast.java @@ -100,7 +100,8 @@ public class ExportFromHazelcast extends Stage implements InitializingBean { writeOk = true; emptyMap = tableRows.isEmpty(); } catch (Exception e) { - log.error(String.format("uuid %s. Can't export table %s to file %s. Table was skipped.", resultContainer.getUuid(), table, dbfFile), e); + log.error(String.format("uuid %s. Can't export table %s (map %s) to file %s. Table was skipped.", + resultContainer.getUuid(), table, table.getHazelcastMapName(), dbfFile), e); return StageResult.ERROR; } finally { if (!writeOk || emptyMap) { diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java index 53dc19b56..3371e1c72 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/DBFExportService.java @@ -1,24 +1,37 @@ package ru.spcex.clearing.dbf.exporter.services; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.stereotype.Service; import ru.spcex.clearing.dbf.exporter.logic.Processor; import ru.spcex.clearing.dbf.exporter.logic.data.ResultContainer; import ru.spcex.clearing.dbf.exporter.logic.data.enums.Table; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.Arrays; @Service("dbfExportService") public class DBFExportService { + Logger log = LoggerFactory.getLogger(getClass()); private final ThreadPoolTaskExecutor executor; private final Processor processor; + protected final ImdgProvider imdgProvider; public DBFExportService(@Qualifier("executor") ThreadPoolTaskExecutor executor, - @Qualifier("processor") Processor processor) { + @Qualifier("processor") Processor processor, + ImdgProvider imdgProvider) { this.executor = executor; this.processor = processor; + this.imdgProvider = imdgProvider; + log.debug("Check IMDG..."); + imdgProvider.waitAvailable(); + log.debug("IMDG ready..."); } public void run() { + log.debug("Do export for all: {}", Arrays.toString(Table.values())); for (Table tableForExport : Table.values()) { executor.submit(() -> processor.process(ResultContainer.createNewTask(tableForExport))); } diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/converters/S_DF08_Converter.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/converters/S_DF08_Converter.java index 62db8cea1..516c4736a 100644 --- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/converters/S_DF08_Converter.java +++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/services/converters/S_DF08_Converter.java @@ -5,6 +5,7 @@ import com.linuxense.javadbf.DBFField; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.sdf.SDf08; +import java.time.format.DateTimeFormatter; import java.util.LinkedList; import java.util.List; @@ -14,7 +15,7 @@ public class S_DF08_Converter extends DFConverter { public Object[] toObjectArray(SDf08 entity) { List values = new LinkedList<>(); values.add(convertStrToLong(entity.getNumber())); - values.add(convertStrToLong(entity.getDatetime())); + values.add(typeMatch(convertStrToLong(entity.getDatetime()))); // UNIX TIME return values.toArray(Object[]::new); } @@ -22,12 +23,16 @@ public class S_DF08_Converter extends DFConverter { public DBFField[] getDBFHeaders() { List dbfFields = new LinkedList<>(); dbfFields.add(new DBFField("NUMBER", DBFDataType.NUMERIC, 10)); - dbfFields.add(new DBFField("DATETIME", DBFDataType.NUMERIC, 13)); + dbfFields.add(new DBFField("DATETIME", DBFDataType.NUMERIC, 13)); // UNIX TIME return dbfFields.toArray(DBFField[]::new); } - private Long convertStrToLong(String s) { + Long convertStrToLong(String s) { if (s == null) return null; return Long.valueOf(s); } + + private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("dd.MM.yyyy"); + private static final DateTimeFormatter DATE_TIME_FMT = DateTimeFormatter.ofPattern("dd.MM.yyyy HH:mm:ss"); + } diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java index eec3f4f88..d79c3a9b9 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/services/LauncherCommandReceiver.java @@ -10,6 +10,7 @@ import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.platform.enumeration.Task; +import ru.spcex.platform.utils.log.ExceptionUtils; @Service public class LauncherCommandReceiver extends QueueConsumer implements InitializingBean { @@ -27,7 +28,20 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi public void afterPropertiesSet() { callback(LauncherCommandRequest.class) .setConsumer(action -> importer.run(ETable.DF_04)) - .forDestination(Task.createOrderConfirm.topic(), callbacks::put); + .forDestination(Task.createOrderConfirm.topic(), callbacks::put); // CORC + callback(LauncherCommandRequest.class) + .setConsumer(action -> importer.run(ETable.DF_01)) + .forDestination(Task.dbf_GBAL.topic(), callbacks::put); // GBAL + callback(LauncherCommandRequest.class) + .setConsumer(action -> importer.run(ETable.DF_12)) + .forDestination(Task.dbf_ABLK.topic(), callbacks::put); // ABLK + callback(LauncherCommandRequest.class) + .setConsumer(action -> importer.run(ETable.DF_09)) + .forDestination(Task.dbf_GBLD.topic(), callbacks::put); // GBLD + callback(LauncherCommandRequest.class) + .setConsumer(action -> importer.run(ETable.DF_16)) + .forDestination(Task.dbf_ADBL.topic(), callbacks::put); // ADBL + init(); } } diff --git a/clearing-parent/reports-service/pom.xml b/clearing-parent/reports-service/pom.xml index 457291ada..969ee84ed 100644 --- a/clearing-parent/reports-service/pom.xml +++ b/clearing-parent/reports-service/pom.xml @@ -24,6 +24,12 @@ org.springframework.boot spring-boot-starter-web + + + org.springframework.boot + spring-boot-starter-tomcat + + org.springframework.boot diff --git a/clearing-parent/reports-service/readme.md b/clearing-parent/reports-service/readme.md index 995cd56a8..4ab9323ef 100644 --- a/clearing-parent/reports-service/readme.md +++ b/clearing-parent/reports-service/readme.md @@ -48,4 +48,9 @@ DAILY - все ежедневные reports-service.kafka-consumer - группа настроек для подключения к очереди kafka -и другие настройки. \ No newline at end of file +и другие настройки. + +Рабочие папки +------------- + +* Из настройки reports-service.ReportModuleOut=./report_out - путь к папке для результатов генерации отчётов diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/notifications/NotificationBuilder.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/notifications/NotificationBuilder.java index 56cdad6b8..9773f4a1e 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/notifications/NotificationBuilder.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/notifications/NotificationBuilder.java @@ -1,19 +1,20 @@ package ru.spcex.clearing.reports.notifications; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import ru.spcex.clearing.reports.exceptions.ConfigException; import ru.spcex.clearing.reports.services.DOCXService; import java.io.File; import java.net.URL; import java.time.LocalDate; -import java.time.LocalDateTime; import java.time.format.DateTimeFormatter; import java.time.format.TextStyle; -import java.util.HashMap; import java.util.Locale; import java.util.Map; public abstract class NotificationBuilder { + protected final Logger log = LoggerFactory.getLogger(getClass()); protected final DateTimeFormatter filenameDateTimeFormatter = DateTimeFormatter.ofPattern("dd.MM.yyyy'T'HH:mm:ss"); protected final DateTimeFormatter dateFormatter = DateTimeFormatter.ofPattern("dd.MM.yyyy"); protected final DOCXService docxService; @@ -41,17 +42,17 @@ public abstract class NotificationBuilder { } - protected File getTemplateFile() { - File templateFile; + protected URL getTemplateFile() { String templateResourceName = getTemplateResourceName(); + URL resource = null; try { - URL resource = getClass().getClassLoader().getResource(templateResourceName); + resource = getClass().getClassLoader().getResource(templateResourceName); if (resource == null) throw new Exception(); - templateFile = new File(resource.toURI()); } catch (Exception e) { - throw new ConfigException("Can't read template resource " + templateResourceName, e); + throw new ConfigException("Can't read template resource " + templateResourceName +" (URL \""+resource+"\")", e); } - return templateFile; + log.info("Use resource: {}", resource); + return resource; } protected String normalizeFilenameForOS(String filename) { diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/DOCXService.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/DOCXService.java index 02e1813eb..67cd092ea 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/DOCXService.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/DOCXService.java @@ -5,6 +5,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.*; +import java.net.URL; import java.util.LinkedList; import java.util.List; import java.util.Map; @@ -43,7 +44,7 @@ public class DOCXService { /** * Файл шаблона */ - private File templateFile; + private URL templateFile; /** * Флаг инициализации @@ -53,13 +54,14 @@ public class DOCXService { /** * (ре)Инициализация сервиса */ - public boolean init(File templateFile, String keyPattern) { + public boolean init(URL templateFile, String keyPattern) { initialized = false; - if (templateFile == null || !templateFile.exists() || !templateFile.isFile()) { - log.warn("Invalid templateFile"); + if (templateFile == null /*|| !templateFile.exists() || !templateFile.isFile()*/) { + log.warn("Invalid templateFile {}", templateFile); return false; } this.templateFile = templateFile; + log.debug("{} use template from {}", getClass().getSimpleName(), templateFile); if (keyPattern != null) { try { this.keyPattern = Pattern.compile(keyPattern); @@ -108,7 +110,7 @@ public class DOCXService { } } - try (InputStream inputStream = new FileInputStream(templateFile); + try (InputStream inputStream = templateFile.openStream(); //new FileInputStream(templateFile); OutputStream outputStream = new FileOutputStream(outputFile)) { XWPFDocument doc = new XWPFDocument(inputStream); @@ -125,7 +127,7 @@ public class DOCXService { boolean checkOk = checkIntersection(valuesMap, notPresentInMap); if (!checkOk) { log.error("Can't insert values in file {}. Keys not in map: {}.", - templateFile.getAbsolutePath(), + templateFile.toExternalForm()/*templateFile.getAbsolutePath()*/, String.join(", ", notPresentInMap)); return false; } diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java index 78adaec08..021bbc884 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/QCommandExecutor.java @@ -42,7 +42,7 @@ public class QCommandExecutor extends QueueConsumer implements InitializingBean log.info("Init queue listener {}", getClass().getSimpleName()); callback(ReportWithPeriodRequest.class) .setConsumer(this::newReportWithPeriod) - .forDestination(Task.createReport.topic(), callbacks::put); + .forDestination(Task.createReport_GREP.topic(), callbacks::put); init(); } diff --git a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/collector/MapUtil.java b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/collector/MapUtil.java index e49ff0ca6..621868c9b 100644 --- a/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/collector/MapUtil.java +++ b/clearing-parent/reports-service/src/main/java/ru/spcex/clearing/reports/services/collector/MapUtil.java @@ -27,7 +27,7 @@ public class MapUtil { K key = indexFieldExtractor.apply(obj); V exist = index.put(key, obj); if (exist != null) { - throw new IllegalStateException(String.format("%s with id=%s,%s has duplicate key=\"{}\"", + throw new IllegalStateException(String.format("%s with id=%s,%s has duplicate key=\"%s\"", exist.getClass().getSimpleName(), exist.getId(), obj.getId(), key )); } diff --git a/clearing-parent/reports-service/src/main/resources/application.properties b/clearing-parent/reports-service/src/main/resources/application.properties index ffa1ebe50..ed2aa2053 100644 --- a/clearing-parent/reports-service/src/main/resources/application.properties +++ b/clearing-parent/reports-service/src/main/resources/application.properties @@ -1,6 +1,3 @@ -server.port=8080 -server.servlet.context-path=/reports_service -spring.main.web-application-type=servlet reports-service.kafka-consumer.bootstrap-servers=localhost:9092 diff --git a/clearing-parent/reports-service/src/test/java/ru/spcex/clearing/reports/services/DOCXServiceTest.java b/clearing-parent/reports-service/src/test/java/ru/spcex/clearing/reports/services/DOCXServiceTest.java index ffd640744..2cfff9e9d 100644 --- a/clearing-parent/reports-service/src/test/java/ru/spcex/clearing/reports/services/DOCXServiceTest.java +++ b/clearing-parent/reports-service/src/test/java/ru/spcex/clearing/reports/services/DOCXServiceTest.java @@ -19,7 +19,7 @@ class DOCXServiceTest { valuesMap.put("${test3}", "======== TEST 3 ========="); DOCXService service = new DOCXService(); - service.init(new File(getClass().getClassLoader().getResource("REC.ACTV_TEST.docx").getPath()), null); + service.init(getClass().getClassLoader().getResource("REC.ACTV_TEST.docx"), null); service.insertValues(new File("output.docx"), valuesMap, false); diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java index 6022af291..8558b06f8 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/SchedulerServiceApplication.java @@ -2,11 +2,29 @@ package ru.spcex.clearing.scheduler; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.ConfigurableApplicationContext; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.PlannerNewRequest; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.scheduler.service.PlannerService; @SpringBootApplication public class SchedulerServiceApplication { public static void main(String[] args) { SpringApplication springApplication = new SpringApplication(SchedulerServiceApplication.class); - springApplication.run(args); + ConfigurableApplicationContext context = springApplication.run(args); + PlannerService plannerService = context.getBean(PlannerService.class); + try { + Thread.sleep(100000); + } catch (InterruptedException e) { + e.printStackTrace(); + } + var baseRequest = new BaseRequest(); + var plannerNewRequest = new PlannerNewRequest(); + plannerNewRequest.task = "GBAL"; + baseRequest.setId(777L); + baseRequest.setRequestPayload(plannerNewRequest); + RequestInfoUpdate requestInfoUpdate = plannerService.newScheduler(baseRequest); + System.out.println(requestInfoUpdate.toString()); } } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ErrorResolverConfig.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ErrorResolverConfig.java new file mode 100644 index 000000000..c64e7c10e --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ErrorResolverConfig.java @@ -0,0 +1,14 @@ +package ru.spcex.clearing.scheduler.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.enumeration.SimpleMessageResolver; + +@Configuration +public class ErrorResolverConfig { + @Bean + public IMessageResolver messageResolver() { + return new SimpleMessageResolver(); + } +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ValidationConfig.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ValidationConfig.java index 3948ed69a..f11175789 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ValidationConfig.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/config/ValidationConfig.java @@ -2,10 +2,6 @@ package ru.spcex.clearing.scheduler.config; 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.company.Company; -import ru.clearing.classes.statics.data.company.CompanySymbols; import ru.clearing.platform.dictionary.TaskDictionary; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.PlannerNewRequest; @@ -30,10 +26,7 @@ public class ValidationConfig { public ValidationConfig(ImdgProvider imdgProvider) { this.imdgs = new HashMap<>(); BiConsumer> addImdg = (s, aClass) -> imdgs.put(s, imdgProvider.getImdg(s, aClass)); - addImdg.accept(IMDGDistributedNames.Map_Account, Account.class); - addImdg.accept(IMDGDistributedNames.Map_Company, Company.class); - addImdg.accept(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); - addImdg.accept(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class); + addImdg.accept(IMDGDistributedNames.Map_TaskDictionary, TaskDictionary.class); } /** @@ -50,7 +43,6 @@ public class ValidationConfig { context.setValidatedObject(plannerNewRequest); Consumer addImdg = (s) -> context.addImdg(s, getImdg(s)); addImdg.accept(IMDGDistributedNames.Map_TaskDictionary); - addImdg.accept(IMDGDistributedNames.Map_TaskStatusDictionary); return new ValidatorImpl<>(context, EnumPresentRule.instance("task", PlannerNewRequest::getTask, diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/error/Errors.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/error/PlannerErrors.java similarity index 68% rename from clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/error/Errors.java rename to clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/error/PlannerErrors.java index 7fdec8a3e..fc15d5a5b 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/error/Errors.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/error/PlannerErrors.java @@ -2,14 +2,13 @@ package ru.spcex.clearing.scheduler.error; import ru.spcex.platform.utils.enumeration.IEnumId; -public enum Errors implements IEnumId { - WrongFieldValue(10003L) - +public enum PlannerErrors implements IEnumId { + WrongEnumValue(10003L), ; private final Long id; - Errors(Long id) { + PlannerErrors(Long id) { this.id = id; } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java new file mode 100644 index 000000000..bbf98c9d5 --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/LauncherSender.java @@ -0,0 +1,68 @@ +package ru.spcex.clearing.scheduler.service; + + +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportWithPeriodRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.platform.enumeration.Task; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; + +@Service +public class LauncherSender { + Logger log = LoggerFactory.getLogger(getClass()); + private final Producer kafka; + private final ImdgId idGenerator; + + @Autowired + public LauncherSender(Producer kafka, ImdgProvider imdgProvider) { + this.kafka = kafka; + this.idGenerator = imdgProvider.getImdgIdGenerator(); + } + + protected BaseRequest makeCmdRequest(Task byTask, Long userId) { + Object toRequest; + if (Task.createReport_GREP.equals(byTask)) { // см. ru.spcex.clearing.backendapi.controller.request.cud.schedule.LauncherNew#toRequest + ReportWithPeriodRequest reportCommand = new ReportWithPeriodRequest(); + reportCommand.setReportId(byTask.getKey()); + toRequest = reportCommand; + } else { + LauncherCommandRequest taskRunnerCommandRequest = new LauncherCommandRequest(); + taskRunnerCommandRequest.setTaskName(byTask.getKey()); + taskRunnerCommandRequest.setUserId(userId); + toRequest = taskRunnerCommandRequest; + } + BaseRequest request = new BaseRequest<>(); + request.setId(idGenerator.nextId()); + request.setActionType(ActionType.NEW); + request.setRequestPayload(toRequest);// iAction.toRequest() + return request; + } + + public void sendCommandToQueue(Task toTaskQueue, Long userId) { + BaseRequest request = makeCmdRequest(toTaskQueue, userId); + String destination = toTaskQueue.topic(); // "launcher-" + getKey() + log.debug("Send command {} to {}", toTaskQueue, destination); + Future send = kafka.send(new ProducerRecord<>(destination, request)); + try { + send.get(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException("Command " + toTaskQueue + " not send", e); + } catch (ExecutionException e) { + throw new RuntimeException("Command " + toTaskQueue + " not send", e); + } + } + +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/NewRequestResult.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/NewRequestResult.java new file mode 100644 index 000000000..282be46cd --- /dev/null +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/NewRequestResult.java @@ -0,0 +1,12 @@ +package ru.spcex.clearing.scheduler.service; + +import ru.spcex.platform.utils.enumeration.EnumMessage; + +public class NewRequestResult { + private EnumMessage error; + + public NewRequestResult(EnumMessage validationError) { + this.error = validationError; + } + +} diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java index 343e57b85..869168f21 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/PlannerService.java @@ -6,6 +6,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; 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.scheduler.Planner; import ru.spcex.clearing.imdg.IMDGDistributedNames; @@ -15,34 +16,48 @@ import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteReques import ru.spcex.clearing.platform.messaging.domain.cud.schedule.PlannerNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.PlannerUpdateRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; import ru.spcex.clearing.scheduler.IRequestValidator; +import ru.spcex.clearing.scheduler.error.PlannerErrors; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.validation.IValidator; import java.time.Instant; import java.util.List; +import java.util.Optional; +import java.util.function.Function; @Service public class PlannerService extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); private final ImdgProvider imdgProvider; private final Imdg plannerMap; - private final IRequestValidator plannerNewRequest = IRequestValidator.PLANNER_NEW_REQUEST; + private final IMessageResolver messageResolver; + private final Function validationFactory; private final IRequestValidator plannerUpdateRequest = IRequestValidator.PLANNER_UPDATE_REQUEST; private final IRequestValidator deleteRequestValidator = IRequestValidator.COMMON_DELETE_REQUEST; @Autowired - public PlannerService(Consumer kafkaQueue, Producer kafkaProducer, - ImdgProvider imdgProvider) { + public PlannerService(Consumer kafkaQueue, + Producer kafkaProducer, + ImdgProvider imdgProvider, + IMessageResolver messageResolver, + @Qualifier("plannerNewRequestValidator") Function plannerNewRequestValidator) { super(kafkaQueue, kafkaProducer); this.plannerMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Planner, Planner.class); + this.messageResolver = messageResolver; + this.validationFactory = plannerNewRequestValidator; this.imdgProvider = imdgProvider; } @Override public void afterPropertiesSet() { callback(PlannerNewRequest.class) - .setConsumer(this::newScheduler) + .setFunction(this::newScheduler) .forDestination(Consts.DESTINATION_PLANNER_NEW, callbacks::put); callback(PlannerUpdateRequest.class) .setConsumer(this::updateScheduler) @@ -53,11 +68,19 @@ public class PlannerService extends QueueConsumer implements InitializingBean { init(); } - private void newScheduler(BaseRequest userRequest) { + public RequestInfoUpdate newScheduler(BaseRequest userRequest) { PlannerNewRequest req = userRequest.getRequestPayload(); log.debug("PlannerNewRequest received"); - List errors = plannerNewRequest.validate(req, imdgProvider); - if (IRequestValidator.checkErrorList(errors, log)) return; + IValidator validator = validationFactory.apply(req); + Optional validationError = validator.tillFirstError(); + if (validationError.isPresent()) { + String errorMsg = messageResolver.resolve(new EnumMessage(PlannerErrors.WrongEnumValue)); + log.error("cannot process PlannerNewRequest id={}: {}", userRequest.getId(), errorMsg); + return new RequestInfoUpdate() + .setId(userRequest.getId()) + .setStatus(Status.Error) + .setMessage(errorMsg); + } Planner planner = new Planner(); planner.setCreated(Instant.now()); planner.setTask(req.getTask()); @@ -69,6 +92,7 @@ public class PlannerService extends QueueConsumer implements InitializingBean { planner.setSecurityId(req.getSecurityId()); plannerMap.insert(planner); log.debug("successfully processed, new id {}", planner.getId()); + return null; } private void updateScheduler(BaseRequest userRequest) { diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java index 89fd9c81c..bcc7fbc07 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/service/TaskManager.java @@ -49,12 +49,15 @@ public class TaskManager implements EntryAddedListener, private Imdg launcherMap; private Imdg plannerAllTodayMap; private ConcurrentHashMap scheduledJobs; + private LauncherSender launcherSender; @Autowired TaskManager(TaskScheduler taskScheduler, - ImdgProvider imdgProvider) { + ImdgProvider imdgProvider, + LauncherSender launcherSender) { this.taskScheduler = taskScheduler; this.imdgProvider = imdgProvider; + this.launcherSender = launcherSender; } private static LocalDateTime dateOldTypeConvert(@NonNull Date oldDate) { @@ -220,8 +223,10 @@ public class TaskManager implements EntryAddedListener, // --- Реализация выполнения задач --- protected void doJob(PlannerAllToday task) { - if (getEnumByKey(Task.class, task.getTask()) != null) { + Task taskE = getEnumByKey(Task.class, task.getTask()); + if (taskE != null) { log.debug("LauncherCommandRequest received"); + // 1. cохранить команду Instant created = Instant.now(); Launcher launcher = new Launcher(); launcher.setTask(task.getTask()); @@ -229,6 +234,8 @@ public class TaskManager implements EntryAddedListener, launcher.setCreated(created); launcher.setUpdated(created); launcherMap.insert(launcher); + // 2. отправить сообщение + launcherSender.sendCommandToQueue(taskE, task.getParentId()); log.debug("successfully processed, new id {}", launcher.getId()); } } diff --git a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/EnumPresentRule.java b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/EnumPresentRule.java index 0acec99af..36eda3401 100644 --- a/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/EnumPresentRule.java +++ b/clearing-parent/scheduler-service/src/main/java/ru/spcex/clearing/scheduler/validation/rules/EnumPresentRule.java @@ -1,6 +1,6 @@ package ru.spcex.clearing.scheduler.validation.rules; -import ru.spcex.clearing.scheduler.error.Errors; +import ru.spcex.clearing.scheduler.error.PlannerErrors; import ru.spcex.platform.classes.base.SpcexObjectBase; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.validation.ImdgValidationContext; @@ -41,7 +41,7 @@ public class EnumPresentRule implements IValidatio ); D taskFromMap = dictImdg.getSingleObjectByFieldValues(Map.of("code", enumCode)); if (taskFromMap == null) { - return of(Errors.WrongFieldValue, fieldName); + return of(PlannerErrors.WrongEnumValue, fieldName); } return empty(); } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java index 1f6ca6675..481aec643 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/ValidationProvider.java @@ -7,6 +7,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSe import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest; import ru.spcex.clearing.securities.errors.SecuritiesError; import ru.spcex.clearing.securities.validation.rule.EndDtAfterStartDt; +import ru.spcex.clearing.securities.validation.rule.MmsDeleteValidationRule; import ru.spcex.clearing.securities.validation.rule.MmsNewValidationRule; import ru.spcex.clearing.securities.validation.rule.MmsUpdateValidationRule; import ru.spcex.platform.classes.base.SpcexObjectBase; @@ -59,7 +60,8 @@ public class ValidationProvider { context.setValidatedObject(mmsRequest); context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, mmsMap); return new ValidatorImpl<>(context, - new PresentById(IMDGDistributedNames.Map_MoneyMarketSecurity, SecuritiesError.InstrumentNotFound, true) + new PresentById(IMDGDistributedNames.Map_MoneyMarketSecurity, SecuritiesError.InstrumentNotFound, true), + MmsDeleteValidationRule.StatusIsActive ); }; } diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsDeleteValidationRule.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsDeleteValidationRule.java new file mode 100644 index 000000000..bbf126857 --- /dev/null +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsDeleteValidationRule.java @@ -0,0 +1,32 @@ +package ru.spcex.clearing.securities.validation.rule; + +import ru.clearing.classes.statics.data.misc.MoneyMarketSecurity; +import ru.spcex.clearing.platform.messaging.domain.cud.securitites.MoneyMarketSecurityUpdateRequest; +import ru.spcex.clearing.securities.errors.SecuritiesError; +import ru.spcex.platform.enumeration.Status; +import ru.spcex.platform.imdg.validation.ImdgValidationContext; +import ru.spcex.platform.imdg.validation.Stored; +import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.validation.IValidationRule; + +import java.util.Optional; + +public enum MmsDeleteValidationRule implements IValidationRule> { + StatusIsActive() { + @Override + public Optional validate(ImdgValidationContext context) { + MoneyMarketSecurity mms = context.getStoredObject(Stored.PresentById); + //проверка только если PresentById найдет объект и сохранит его + if (mms != null && !Status.Active.getKey().equals(mms.getWorkflowStatus())) { + return of(SecuritiesError.InstrumentNotActive); + } else { + return empty(); + } + } + }; + + @Override + public String ruleName() { + return "MmsDeleteValidationRule." + name(); + } +} diff --git a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsUpdateValidationRule.java b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsUpdateValidationRule.java index 2bf92e6c4..1d88b45bb 100644 --- a/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsUpdateValidationRule.java +++ b/clearing-parent/securities-service/src/main/java/ru/spcex/clearing/securities/validation/rule/MmsUpdateValidationRule.java @@ -27,6 +27,6 @@ public enum MmsUpdateValidationRule implements IValidationRule