for test
This commit is contained in:
commit
5167eb45f0
43 changed files with 562 additions and 92 deletions
|
|
@ -24,6 +24,10 @@
|
|||
<groupId>ru.spcex.platform</groupId>
|
||||
<artifactId>platform-imdg-api-hazelcast-impl</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ru.spcex.platform</groupId>
|
||||
<artifactId>platform-enum</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ru.spcex.clearing</groupId>
|
||||
<artifactId>classes</artifactId>
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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<BankAccount> bankAccountMap;
|
||||
private final Imdg<Account> accountMap;
|
||||
private final Imdg<Relation> relationMap;
|
||||
private final IMessageResolver messageResolver;
|
||||
|
||||
@Autowired
|
||||
public BankAccountService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
|
||||
public BankAccountService(Consumer<String, Object> kafkaQueue,
|
||||
Producer<String, Object> 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<BankAccountNewRequest> userRequest) {
|
||||
private RequestInfoUpdate bankAccountNew(BaseRequest<BankAccountNewRequest> userRequest) {
|
||||
BankAccountNewRequest req = userRequest.getRequestPayload();
|
||||
|
||||
Collection<Account> 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<BankAccountUpdateRequest> userRequest) {
|
||||
private RequestInfoUpdate bankAccountUpdate(BaseRequest<BankAccountUpdateRequest> 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<CommonDeleteRequest> userRequest) {
|
||||
private RequestInfoUpdate bankAccountDelete(BaseRequest<CommonDeleteRequest> 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;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<BankAccount> 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<String, Object> mockConsumer;
|
||||
private MockProducer<String, Object> 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();
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<LauncherCommandRequest> {
|
||||
public class LauncherNew implements IAction<Object> {
|
||||
@ApiModelProperty(value = "Идентификатор единоличного исполнительного органа", example = "ABCD")
|
||||
@JsonProperty
|
||||
private String task;
|
||||
|
|
@ -22,11 +24,17 @@ public class LauncherNew implements IAction<LauncherCommandRequest> {
|
|||
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
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -137,6 +137,9 @@ public class Sdf09Executor extends AbstractExecutor<SDf09> {
|
|||
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);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -58,7 +58,6 @@ public enum Sdf01ValidationRule implements IValidationRule<ImdgValidationContext
|
|||
@Override
|
||||
public Optional<EnumMessage> validate(ImdgValidationContext<SDf01> context) {
|
||||
SDf01 sdf01 = context.getValidatedObject();
|
||||
//fixme string format???
|
||||
if (!LocalDate.now().equals(LocalDate.parse(sdf01.getDat(), datFormatter))) {
|
||||
return of(BalanceError.CurrentDateOnly);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<ImdgValidationContext
|
|||
SDf09 sdf09 = context.getValidatedObject();
|
||||
Imdg<CompanySymbols> companySymbolImdg = context.obtainMap(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
|
||||
Imdg<Company> 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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<ImdgValidationContext<SDf16>> {
|
||||
|
|
@ -20,11 +23,26 @@ public enum Sdf16ValidationRule implements IValidationRule<ImdgValidationContext
|
|||
SDf16 sdf16 = context.getValidatedObject();
|
||||
Imdg<CompanySymbols> companySymbolsImdg = context.obtainMap(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class);
|
||||
Imdg<Company> 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());
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<SDf16> 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<EnumMessage> err = Sdf16ValidationRule.CompanyPresent.validate(ctx);
|
||||
assertTrue(err.isPresent());
|
||||
}
|
||||
|
||||
<T extends SpcexObjectBase> Imdg<T> makeTestMap() {
|
||||
return new Imdg<T>() {
|
||||
@Override
|
||||
public String getMapName() {
|
||||
return Imdg.super.getMapName();
|
||||
}
|
||||
|
||||
@Override
|
||||
public T getSingleObjectBySQL(String sql) {
|
||||
return null;
|
||||
}
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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)));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<SDf08> {
|
|||
public Object[] toObjectArray(SDf08 entity) {
|
||||
List<Object> 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<SDf08> {
|
|||
public DBFField[] getDBFHeaders() {
|
||||
List<DBFField> 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");
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -24,6 +24,12 @@
|
|||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-web</artifactId>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-tomcat</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
|
|
|
|||
|
|
@ -48,4 +48,9 @@ DAILY - все ежедневные
|
|||
|
||||
reports-service.kafka-consumer - группа настроек для подключения к очереди kafka
|
||||
|
||||
и другие настройки.
|
||||
и другие настройки.
|
||||
|
||||
Рабочие папки
|
||||
-------------
|
||||
|
||||
* Из настройки reports-service.ReportModuleOut=./report_out - путь к папке для результатов генерации отчётов
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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<PlannerNewRequest>();
|
||||
var plannerNewRequest = new PlannerNewRequest();
|
||||
plannerNewRequest.task = "GBAL";
|
||||
baseRequest.setId(777L);
|
||||
baseRequest.setRequestPayload(plannerNewRequest);
|
||||
RequestInfoUpdate requestInfoUpdate = plannerService.newScheduler(baseRequest);
|
||||
System.out.println(requestInfoUpdate.toString());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String, Class<? extends SpcexObjectBase>> 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<String> 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,
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
||||
|
|
@ -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<String, Object> kafka;
|
||||
private final ImdgId idGenerator;
|
||||
|
||||
@Autowired
|
||||
public LauncherSender(Producer<String, Object> kafka, ImdgProvider imdgProvider) {
|
||||
this.kafka = kafka;
|
||||
this.idGenerator = imdgProvider.getImdgIdGenerator();
|
||||
}
|
||||
|
||||
protected BaseRequest<Object> 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<Object> 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<Object> request = makeCmdRequest(toTaskQueue, userId);
|
||||
String destination = toTaskQueue.topic(); // "launcher-" + getKey()
|
||||
log.debug("Send command {} to {}", toTaskQueue, destination);
|
||||
Future<RecordMetadata> 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);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -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<Planner> plannerMap;
|
||||
private final IRequestValidator<PlannerNewRequest> plannerNewRequest = IRequestValidator.PLANNER_NEW_REQUEST;
|
||||
private final IMessageResolver messageResolver;
|
||||
private final Function<PlannerNewRequest, IValidator> validationFactory;
|
||||
private final IRequestValidator<PlannerUpdateRequest> plannerUpdateRequest = IRequestValidator.PLANNER_UPDATE_REQUEST;
|
||||
private final IRequestValidator<CommonDeleteRequest> deleteRequestValidator = IRequestValidator.COMMON_DELETE_REQUEST;
|
||||
|
||||
@Autowired
|
||||
public PlannerService(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaProducer,
|
||||
ImdgProvider imdgProvider) {
|
||||
public PlannerService(Consumer<String, Object> kafkaQueue,
|
||||
Producer<String, Object> kafkaProducer,
|
||||
ImdgProvider imdgProvider,
|
||||
IMessageResolver messageResolver,
|
||||
@Qualifier("plannerNewRequestValidator") Function<PlannerNewRequest, IValidator> 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<PlannerNewRequest> userRequest) {
|
||||
public RequestInfoUpdate newScheduler(BaseRequest<PlannerNewRequest> userRequest) {
|
||||
PlannerNewRequest req = userRequest.getRequestPayload();
|
||||
log.debug("PlannerNewRequest received");
|
||||
List<IRequestValidator.ValidationError> errors = plannerNewRequest.validate(req, imdgProvider);
|
||||
if (IRequestValidator.checkErrorList(errors, log)) return;
|
||||
IValidator validator = validationFactory.apply(req);
|
||||
Optional<EnumMessage> 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<PlannerUpdateRequest> userRequest) {
|
||||
|
|
|
|||
|
|
@ -49,12 +49,15 @@ public class TaskManager implements EntryAddedListener<Long, PlannerAllToday>,
|
|||
private Imdg<Launcher> launcherMap;
|
||||
private Imdg<PlannerAllToday> plannerAllTodayMap;
|
||||
private ConcurrentHashMap<LocalTime, ScheduledFuture> 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<Long, PlannerAllToday>,
|
|||
// --- Реализация выполнения задач ---
|
||||
|
||||
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<Long, PlannerAllToday>,
|
|||
launcher.setCreated(created);
|
||||
launcher.setUpdated(created);
|
||||
launcherMap.insert(launcher);
|
||||
// 2. отправить сообщение
|
||||
launcherSender.sendCommandToQueue(taskE, task.getParentId());
|
||||
log.debug("successfully processed, new id {}", launcher.getId());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<T, D extends SpcexObjectBase> 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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
);
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<ImdgValidationContext<MoneyMarketSecurityUpdateRequest>> {
|
||||
StatusIsActive() {
|
||||
@Override
|
||||
public Optional<EnumMessage> validate(ImdgValidationContext<MoneyMarketSecurityUpdateRequest> 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();
|
||||
}
|
||||
}
|
||||
|
|
@ -27,6 +27,6 @@ public enum MmsUpdateValidationRule implements IValidationRule<ImdgValidationCon
|
|||
|
||||
@Override
|
||||
public String ruleName() {
|
||||
return "MmsNewValidationRule." + name();
|
||||
return "MmsUpdateValidationRule." + name();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,18 @@
|
|||
package ru.spcex.platform.enumeration;
|
||||
|
||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum AccountStatus implements IEnumKey {
|
||||
ACTIVE("ACTV"), BLOCKED("BLKD"), CLOSE("CLOS");
|
||||
|
||||
private final String key;
|
||||
|
||||
AccountStatus(String key) {
|
||||
this.key = key;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getKey() {
|
||||
return key;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,18 @@
|
|||
package ru.spcex.platform.enumeration;
|
||||
|
||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum Allowed implements IEnumKey {
|
||||
ALLOWED("ALWD"), DENIED("DEND");
|
||||
|
||||
private final String key;
|
||||
|
||||
Allowed(String key) {
|
||||
this.key = key;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getKey() {
|
||||
return key;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,18 @@
|
|||
package ru.spcex.platform.enumeration;
|
||||
|
||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum Service implements IEnumKey {
|
||||
MKR("MKR");
|
||||
|
||||
private final String key;
|
||||
|
||||
Service(String key) {
|
||||
this.key = key;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getKey() {
|
||||
return key;
|
||||
}
|
||||
}
|
||||
|
|
@ -13,10 +13,13 @@ public enum Task implements IEnumKey {
|
|||
startOfPreClearing("SPRC"),// Запуск преклиринга
|
||||
startPostClearing("SPOC"),// Запуск постклиринга
|
||||
createOrder("CORD"),
|
||||
dbf_GBAL("GBAL"),
|
||||
dbf_ABLK("ABLK"),
|
||||
dbf_GBLD("GBLD"),
|
||||
dbf_ADBL("ADBL"),
|
||||
createOrderConfirm("CORC"),
|
||||
getAllBalance("GALB"),
|
||||
createReport("RPRT");// Создание отчёта (report-service)
|
||||
|
||||
createReport_GREP("GREP");// Создание отчёта (report-service) RPRT нескольких видов, этот GREP
|
||||
|
||||
private final String key;
|
||||
|
||||
|
|
|
|||
|
|
@ -4,23 +4,25 @@ import com.fasterxml.jackson.annotation.JsonProperty;
|
|||
|
||||
public class BankAccountNewRequest {
|
||||
@JsonProperty
|
||||
public String bankIdentificationCode;//Банковский идентификационный код (БИК)
|
||||
public String bankIdentificationCode;
|
||||
@JsonProperty
|
||||
public String bankName;//Наименование банка
|
||||
public String bankName;
|
||||
@JsonProperty
|
||||
public String correspondentAccount;//Корреспондентский счет
|
||||
public String correspondentAccount;
|
||||
@JsonProperty
|
||||
public String correspondentAccountName;//Наименование корреспондентского счета
|
||||
public String correspondentAccountName;
|
||||
@JsonProperty
|
||||
public String currency;//Идентификатор валюты
|
||||
public String currency;
|
||||
@JsonProperty
|
||||
public String destination;//Назначение
|
||||
public String destination;
|
||||
@JsonProperty
|
||||
public String taxpayerIdentificationNumber;//Идентификационный номер налогоплательщика (ИНН)
|
||||
public String taxpayerIdentificationNumber;
|
||||
@JsonProperty
|
||||
public String taxRegistrationReasonCode;//Код причины постановки (КПП)
|
||||
public String taxRegistrationReasonCode;
|
||||
@JsonProperty
|
||||
public String account;//Номер счета
|
||||
public String account;
|
||||
@JsonProperty
|
||||
public Long companyId;
|
||||
|
||||
public String getBankIdentificationCode() {
|
||||
return bankIdentificationCode;
|
||||
|
|
|
|||
|
|
@ -39,4 +39,13 @@ public class RequestInfoUpdate {
|
|||
this.message = message;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "RequestInfoUpdate{" +
|
||||
"id=" + id +
|
||||
", status=" + status +
|
||||
", message='" + message + '\'' +
|
||||
'}';
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue