This commit is contained in:
parent
4ad0fc2a94
commit
db38a4230a
16 changed files with 454 additions and 65 deletions
|
|
@ -1,18 +1,46 @@
|
||||||
package ru.spcex.clearing.account.config;
|
package ru.spcex.clearing.account.config;
|
||||||
|
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.context.annotation.Bean;
|
import org.springframework.context.annotation.Bean;
|
||||||
import org.springframework.context.annotation.Configuration;
|
import org.springframework.context.annotation.Configuration;
|
||||||
import ru.spcex.clearing.account.config.settings.AccountServiceSettings;
|
import ru.spcex.clearing.account.config.settings.AccountServiceSettings;
|
||||||
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
|
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgId;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
|
||||||
@Configuration
|
@Configuration
|
||||||
public class KafkaConfig {
|
public class KafkaConfig {
|
||||||
@Autowired
|
@Autowired
|
||||||
@Bean
|
@Bean
|
||||||
public Consumer<String, Object> createProducer(AccountServiceSettings settings) {
|
public Consumer<String, Object> createConsumer(AccountServiceSettings settings) {
|
||||||
return KafkaConsumerFactory.consumer(settings.getKafka());
|
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Bean
|
||||||
|
public Producer<String, Object> createProducer(AccountServiceSettings settings) {
|
||||||
|
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Bean
|
||||||
|
public KafkaSender kafkaSender(Producer<String, Object> kafkaProducer, ImdgProvider imdgProvider) {
|
||||||
|
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
|
||||||
|
return KafkaSender
|
||||||
|
.setup()
|
||||||
|
.producer(kafkaProducer)
|
||||||
|
.idGenerator(imdgIdGenerator::nextId)
|
||||||
|
.imdgProvider(s -> {
|
||||||
|
Imdg<RequestInfo> imdg = imdgProvider.getImdg(s, RequestInfo.class);
|
||||||
|
return imdg::insert;
|
||||||
|
})
|
||||||
|
.build();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||||
import org.springframework.context.annotation.PropertySource;
|
import org.springframework.context.annotation.PropertySource;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
|
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
|
||||||
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
|
import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
|
||||||
|
|
||||||
@Component
|
@Component
|
||||||
|
|
@ -11,7 +12,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
|
||||||
@ConfigurationProperties("account-service")
|
@ConfigurationProperties("account-service")
|
||||||
public class AccountServiceSettings {
|
public class AccountServiceSettings {
|
||||||
private HazelcastClientParams hazelcast;
|
private HazelcastClientParams hazelcast;
|
||||||
private KafkaConsumerSettings kafka;
|
private KafkaConsumerSettings kafkaConsumer;
|
||||||
|
private KafkaProducerSettings kafkaProducer;
|
||||||
|
|
||||||
public HazelcastClientParams getHazelcast() {
|
public HazelcastClientParams getHazelcast() {
|
||||||
return hazelcast;
|
return hazelcast;
|
||||||
|
|
@ -21,11 +23,19 @@ public class AccountServiceSettings {
|
||||||
this.hazelcast = hazelcast;
|
this.hazelcast = hazelcast;
|
||||||
}
|
}
|
||||||
|
|
||||||
public KafkaConsumerSettings getKafka() {
|
public KafkaConsumerSettings getKafkaConsumer() {
|
||||||
return kafka;
|
return kafkaConsumer;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setKafka(KafkaConsumerSettings kafka) {
|
public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) {
|
||||||
this.kafka = kafka;
|
this.kafkaConsumer = kafkaConsumer;
|
||||||
|
}
|
||||||
|
|
||||||
|
public KafkaProducerSettings getKafkaProducer() {
|
||||||
|
return kafkaProducer;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setKafkaProducer(KafkaProducerSettings kafkaProducer) {
|
||||||
|
this.kafkaProducer = kafkaProducer;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,75 @@
|
||||||
|
package ru.spcex.clearing.account.service;
|
||||||
|
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
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.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01RequestPart;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountSdf01ToStatementRequestPart;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class AccountService extends QueueConsumer implements InitializingBean {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final Imdg<Account> accountMap;
|
||||||
|
private final KafkaSender kafkaSender;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public AccountService(Consumer<String, Object> kafkaQueue, ImdgProvider imdgProvider, KafkaSender kafkaSender) {
|
||||||
|
super(kafkaQueue);
|
||||||
|
this.accountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
||||||
|
this.kafkaSender = kafkaSender;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() {
|
||||||
|
callback(AccountSdf01Request.class)
|
||||||
|
.setConsumer(this::accountNew)
|
||||||
|
.forDestination(Consts.ACCOUNT_NEW, callbacks::put);
|
||||||
|
init();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void accountNew(BaseRequest<AccountSdf01Request> userRequest) {
|
||||||
|
AccountSdf01Request req = userRequest.getRequestPayload();
|
||||||
|
log.debug("AccountSdf01Request received");
|
||||||
|
List<AccountSdf01ToStatementRequestPart> accountToStatement = new ArrayList<>();
|
||||||
|
for (AccountSdf01RequestPart accountReq : req.getAccounts()) {
|
||||||
|
Account account = new Account();
|
||||||
|
account.setAccount(accountReq.getAccount());
|
||||||
|
//fixme account.setCompany();
|
||||||
|
accountMap.insert(account);
|
||||||
|
AccountSdf01ToStatementRequestPart responsePart = responsePart(accountReq.getSdf01Id());
|
||||||
|
accountToStatement.add(responsePart);
|
||||||
|
}
|
||||||
|
sendStatementRequestBack(accountToStatement);
|
||||||
|
log.debug("successfully processed, grouping id={}, processed number={}", req.getGroupingSdf01Id(), accountToStatement.size());
|
||||||
|
}
|
||||||
|
|
||||||
|
private AccountSdf01ToStatementRequestPart responsePart(Long sdf01Id) {
|
||||||
|
AccountSdf01ToStatementRequestPart responsePart = new AccountSdf01ToStatementRequestPart();
|
||||||
|
responsePart.setSdf01Id(sdf01Id);
|
||||||
|
responsePart.setErrorCode(null);
|
||||||
|
responsePart.setErrorText(null);
|
||||||
|
return responsePart;
|
||||||
|
}
|
||||||
|
|
||||||
|
private void sendStatementRequestBack(List<AccountSdf01ToStatementRequestPart> results) {
|
||||||
|
StatementRequest request = new StatementRequest();
|
||||||
|
request.setAccountCreationResults(results);
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -4,6 +4,7 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.context.annotation.Bean;
|
import org.springframework.context.annotation.Bean;
|
||||||
import org.springframework.context.annotation.Configuration;
|
import org.springframework.context.annotation.Configuration;
|
||||||
import ru.clearing.classes.statics.data.account.Account;
|
import ru.clearing.classes.statics.data.account.Account;
|
||||||
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf01;
|
import ru.clearing.classes.statics.data.sdf.SDf01;
|
||||||
import ru.spcex.clearing.balance.validation.Sdf01ValidationRule;
|
import ru.spcex.clearing.balance.validation.Sdf01ValidationRule;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
|
@ -26,7 +27,9 @@ public class ValidationConfig {
|
||||||
context.setValidatedObject(sDf01);
|
context.setValidatedObject(sDf01);
|
||||||
BiConsumer<String, Class<? extends SpcexObjectBase>> addImdg = (s, aClass) -> context.addImdg(s, imdgProvider.getImdg(s, aClass));
|
BiConsumer<String, Class<? extends SpcexObjectBase>> addImdg = (s, aClass) -> context.addImdg(s, imdgProvider.getImdg(s, aClass));
|
||||||
addImdg.accept(IMDGDistributedNames.Map_Account, Account.class);
|
addImdg.accept(IMDGDistributedNames.Map_Account, Account.class);
|
||||||
|
addImdg.accept(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
return new ValidatorImpl<>(context,
|
return new ValidatorImpl<>(context,
|
||||||
|
Sdf01ValidationRule.CompanyPresent,
|
||||||
Sdf01ValidationRule.AccountPresent,
|
Sdf01ValidationRule.AccountPresent,
|
||||||
Sdf01ValidationRule.CurrencyCode,
|
Sdf01ValidationRule.CurrencyCode,
|
||||||
Sdf01ValidationRule.CurrentDateOnly,
|
Sdf01ValidationRule.CurrentDateOnly,
|
||||||
|
|
|
||||||
|
|
@ -10,35 +10,46 @@ import org.springframework.stereotype.Service;
|
||||||
import ru.clearing.classes.statics.data.account.Account;
|
import ru.clearing.classes.statics.data.account.Account;
|
||||||
import ru.clearing.classes.statics.data.company.Company;
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf01;
|
import ru.clearing.classes.statics.data.sdf.SDf01;
|
||||||
|
import ru.clearing.classes.statics.data.sdf.SDf02;
|
||||||
import ru.clearing.classes.statics.data.statement.Statement;
|
import ru.clearing.classes.statics.data.statement.Statement;
|
||||||
import ru.spcex.clearing.balance.errors.BalanceError;
|
import ru.spcex.clearing.balance.errors.BalanceError;
|
||||||
import ru.spcex.clearing.balance.validation.ValidationStored;
|
import ru.spcex.clearing.balance.validation.ValidationStored;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.AccountNewRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01RequestPart;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
|
import ru.spcex.platform.enumeration.CurrencyCode;
|
||||||
|
import ru.spcex.platform.enumeration.InOutSDfType;
|
||||||
|
import ru.spcex.platform.enumeration.OperationStatus;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
|
import ru.spcex.platform.utils.time.TimeUtil;
|
||||||
import ru.spcex.platform.utils.validation.IValidator;
|
import ru.spcex.platform.utils.validation.IValidator;
|
||||||
|
|
||||||
import java.util.Map;
|
import java.time.Instant;
|
||||||
import java.util.Optional;
|
import java.time.LocalDate;
|
||||||
|
import java.time.format.DateTimeFormatter;
|
||||||
|
import java.util.*;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class StatementService extends QueueConsumer implements InitializingBean {
|
public class StatementService extends QueueConsumer implements InitializingBean {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
|
private final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy");
|
||||||
private final ImdgProvider imdgProvider;
|
private final ImdgProvider imdgProvider;
|
||||||
private final KafkaSender kafkaReqProducer;
|
private final KafkaSender kafkaReqProducer;
|
||||||
private final LoggingService errorLogger;
|
private final LoggingService errorLogger;
|
||||||
private final Imdg<SDf01> sdf01Imdg;
|
private final Imdg<SDf01> sdf01Imdg;
|
||||||
private final Imdg<Company> companyImdg;
|
private final Imdg<SDf02> sdf02Imdg;
|
||||||
private final Imdg<Account> accountImdg;
|
private final Imdg<AccountBalance> accountBalanceImdg;
|
||||||
private final Imdg<Statement> statementImdg;
|
private final Imdg<Statement> statementImdg;
|
||||||
private final Function<SDf01, IValidator> sDf01Validator;
|
private final Function<SDf01, IValidator> sDf01Validator;
|
||||||
|
|
||||||
|
|
@ -48,9 +59,9 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
||||||
super(kafkaQueue);
|
super(kafkaQueue);
|
||||||
this.imdgProvider = imdgProvider;
|
this.imdgProvider = imdgProvider;
|
||||||
this.sdf01Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
|
this.sdf01Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf01, SDf01.class);
|
||||||
this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
|
this.sdf02Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf02, SDf02.class);;
|
||||||
this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
|
|
||||||
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
|
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
|
||||||
|
this.accountBalanceImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_AccountBalance, AccountBalance.class);
|
||||||
this.kafkaReqProducer = kafkaReqProducer;
|
this.kafkaReqProducer = kafkaReqProducer;
|
||||||
this.errorLogger = errorLogger;
|
this.errorLogger = errorLogger;
|
||||||
this.sDf01Validator = sDf01Validator;
|
this.sDf01Validator = sDf01Validator;
|
||||||
|
|
@ -60,46 +71,170 @@ public class StatementService extends QueueConsumer implements InitializingBean
|
||||||
public void afterPropertiesSet() throws Exception {
|
public void afterPropertiesSet() throws Exception {
|
||||||
callback(StatementRequest.class)
|
callback(StatementRequest.class)
|
||||||
.setConsumer(this::process)
|
.setConsumer(this::process)
|
||||||
.forDestination(Consts.DESTINATION_SDF02_NEW, callbacks::put);
|
.forDestination(Consts.STATEMENT_PROCESS, callbacks::put);
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
private void process(BaseRequest<StatementRequest> systemRequest) {
|
private void process(BaseRequest<StatementRequest> systemRequest) {
|
||||||
StatementRequest sdfInfo = systemRequest.getRequestPayload();
|
StatementRequest statementRequest = systemRequest.getRequestPayload();
|
||||||
SDf01 sdf01 = sdf01Imdg.getSingleObjectByID(sdfInfo.getSdf01Id());
|
Collection<SDf01> sdf01Group;
|
||||||
Company company = companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", sdf01.getDeal()));
|
if (statementRequest.getAccountCreationResults().size() == 0) {
|
||||||
if (company == null) {
|
sdf01Group = sdf01Imdg.getCollectionObjectsByFieldValues(Map.of("generationId", statementRequest.getSdf01GroupId()));
|
||||||
errorLogger.logError("sdf01.id={}", new EnumMessage(BalanceError.CompanyNotFound), sdfInfo.getSdf01Id());
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
IValidator validator = sDf01Validator.apply(sdf01);
|
|
||||||
Optional<EnumMessage> error = validator.tillFirstError();
|
|
||||||
if (BalanceError.AccountNotPresent.equals(error.map(EnumMessage::getSubject).orElse(null))) {
|
|
||||||
kafkaReqProducer.sendRequestToQueue(Consts.ACCOUNT_NEW, createAccountRequest(sdf01.getAccount()));
|
|
||||||
log.info("account not found - send request for creation");
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
if (error.isPresent()) {
|
|
||||||
errorLogger.logError("sdf01.id={}", error.get(), sdf01.getId());
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
Statement statement = statementImdg.getSingleObjectByFieldValues(Map.of("account", sdf01.getAccount()));
|
|
||||||
|
|
||||||
if (statement == null) {
|
|
||||||
createFlow(sdfInfo, validator.getStored(ValidationStored.Sdf01Account));
|
|
||||||
} else {
|
} else {
|
||||||
// updateFlow(statement);
|
sdf01Group = statementRequest.getAccountCreationResults()
|
||||||
|
.stream()
|
||||||
|
.filter(part -> part.getErrorCode() == null) //fixme эти случае должны попадать в ошибочный sdf02
|
||||||
|
.map(part -> sdf01Imdg.getSingleObjectByID(part.getSdf01Id()))
|
||||||
|
.collect(Collectors.toList());
|
||||||
}
|
}
|
||||||
|
sdf01Group = sdf01Group
|
||||||
|
.stream()
|
||||||
|
.sorted(Comparator.comparing(SpcexObjectBase::getId))
|
||||||
|
.collect(Collectors.toList());
|
||||||
|
List<AccountSdf01RequestPart> accountRequests = new ArrayList<>();
|
||||||
|
for (SDf01 sdf01 : sdf01Group) {
|
||||||
|
IValidator validator = sDf01Validator.apply(sdf01);
|
||||||
|
Optional<EnumMessage> error = validator.tillFirstError();
|
||||||
|
Company company = validator.getStored(ValidationStored.Sdf01Company);
|
||||||
|
if (statementRequest.getAccountCreationResults().size() == 0
|
||||||
|
&& BalanceError.AccountNotPresent.equals(error.map(EnumMessage::getSubject).orElse(null))) {
|
||||||
|
//на данном шаге company существует -> getId ok
|
||||||
|
//формируем пакетный запрос на добавление account
|
||||||
|
//ответ придет в этот же метод, process
|
||||||
|
accountRequests.add(createAccountRequestPart(sdf01.getId(), sdf01.getAccount(), company.getId()));
|
||||||
|
log.info("account {} for sdf01.id={} not found - send request for creation", sdf01.getAccount(), sdf01.getId());
|
||||||
|
} else if (BalanceError.AccountNotPresent.equals(error.map(EnumMessage::getSubject).orElse(null))) {
|
||||||
|
log.error("fatal error: resumed processing after generating accounts, but no account found for sdf01.id={}", sdf01.getId());
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (error.isPresent()) {
|
||||||
|
errorLogger.logError("sdf01.id={}", error.get(), sdf01.getId());
|
||||||
|
sdf02Imdg.insert(createErrorSdf02(sdf01, error.get()));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
Statement statement = statementImdg.getSingleObjectByFieldValues(Map.of("account", sdf01.getAccount()));
|
||||||
|
if (statement == null) {
|
||||||
|
statement = createFlow(sdf01,
|
||||||
|
company,
|
||||||
|
validator.getStored(ValidationStored.Sdf01Account));
|
||||||
|
} else {
|
||||||
|
updateFlow(statement, sdf01,
|
||||||
|
validator.getStored(ValidationStored.Sdf01Account));
|
||||||
|
}
|
||||||
|
SDf02 sdf02New = createSuccessSdf02(sdf01);
|
||||||
|
sdf02Imdg.insert(sdf02New);
|
||||||
|
statement.setOutSDfId(sdf02New.getId());
|
||||||
|
AccountBalance accountBalance = createSuccessAccountBalance();
|
||||||
|
accountBalanceImdg.insert(accountBalance);
|
||||||
|
//fixme possible error from accountBalanceCreation -> other status + errorCode, errorText
|
||||||
|
statement.setStatus(OperationStatus.Executed.getKey());
|
||||||
|
statementImdg.update(statement);
|
||||||
|
|
||||||
|
if (accountRequests.size() > 0) {
|
||||||
|
kafkaReqProducer.sendRequestToQueue(Consts.ACCOUNT_NEW, createAccountsRequest(sdf01.getGenerationId(), accountRequests));
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void createFlow(StatementRequest sdfInfo, Account storedObject) {
|
private Statement createFlow(SDf01 sdf01, Company company, Account account) {
|
||||||
|
Statement statement = new Statement();
|
||||||
|
statement.setAddresseeId(company.getId());
|
||||||
|
statement.setSenderId(2L);// fixme справочник
|
||||||
|
statement.setCreated(Instant.now());
|
||||||
|
statement.setClearingDate(TimeUtil.today());
|
||||||
|
statement.setStatementTypeId(0L); //fixme справочник
|
||||||
|
statement.setAccountId(account.getId());
|
||||||
|
statement.setAccount(0L); //fixme
|
||||||
|
statement.setInOutDirection(0L); //fixme
|
||||||
|
statement.setSettlementDate(TimeUtil.localDateToInstant(LocalDate.parse(sdf01.getDat(), datFormatter)));
|
||||||
|
statement.setAmount(null); //fixme sdf01.getRemainder()
|
||||||
|
statement.setCashMovementCurrencyCode(CurrencyCode.RUB.getKey());
|
||||||
|
statement.setStatus(OperationStatus.Pending.getKey());
|
||||||
|
statement.setInSDfId(sdf01.getId());
|
||||||
|
statement.setInOutSDfType(InOutSDfType.type1.getKey());
|
||||||
|
statementImdg.insert(statement);
|
||||||
|
return statement;
|
||||||
}
|
}
|
||||||
|
|
||||||
private AccountNewRequest createAccountRequest(String account) {
|
private void updateFlow(Statement statement, SDf01 sdf01, Account account) {
|
||||||
AccountNewRequest req = new AccountNewRequest();
|
statement.setUpdated(Instant.now());
|
||||||
|
statement.setStatementTypeId(0L); //fixme справочник
|
||||||
|
statement.setAccountId(account.getId());
|
||||||
|
statement.setAccount(0L); //fixme
|
||||||
|
statement.setInOutDirection(0L); //fixme
|
||||||
|
statement.setSettlementDate(TimeUtil.localDateToInstant(LocalDate.parse(sdf01.getDat(), datFormatter)));
|
||||||
|
statement.setAmount(null); //fixme sdf01.getRemainder()
|
||||||
|
statement.setCashMovementCurrencyCode(CurrencyCode.RUB.getKey());
|
||||||
|
statementImdg.update(statement);
|
||||||
|
}
|
||||||
|
|
||||||
|
private AccountSdf01RequestPart createAccountRequestPart(Long sdf01Id, String account, Long companyId) {
|
||||||
|
AccountSdf01RequestPart req = new AccountSdf01RequestPart();
|
||||||
req.setAccount(account);
|
req.setAccount(account);
|
||||||
|
req.setCompanyId(companyId);
|
||||||
|
req.setSdf01Id(sdf01Id);
|
||||||
return req;
|
return req;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private AccountSdf01Request createAccountsRequest(Long sdf01GroupingId, List<AccountSdf01RequestPart> accountRequests) {
|
||||||
|
AccountSdf01Request r = new AccountSdf01Request();
|
||||||
|
r.setGroupingSdf01Id(sdf01GroupingId);
|
||||||
|
r.setAccounts(accountRequests);
|
||||||
|
return r;
|
||||||
|
}
|
||||||
|
|
||||||
|
private SDf02 createErrorSdf02(SDf01 sdf01, EnumMessage error) {
|
||||||
|
SDf02 sDf02 = new SDf02();
|
||||||
|
sDf02.setCurr_code(sdf01.getCurr_code());
|
||||||
|
sDf02.setAccount(sdf01.getAccount());
|
||||||
|
sDf02.setRemainder(sdf01.getRemainder());
|
||||||
|
sDf02.setDeal(sdf01.getDeal());
|
||||||
|
sDf02.setAcc_code(sdf01.getAcc_code());
|
||||||
|
sDf02.setDat(sdf01.getDat());
|
||||||
|
sDf02.setMarket(sdf01.getMarket());
|
||||||
|
sDf02.setAcc_name(sdf01.getAcc_name());
|
||||||
|
sDf02.setAcc_type(sdf01.getAcc_type());
|
||||||
|
sDf02.setSumengage(sdf01.getSumengage());
|
||||||
|
sDf02.setSumunblock(sdf01.getSumunblock());
|
||||||
|
sDf02.setFile_type(sdf01.getFile_type());
|
||||||
|
sDf02.setInSDf01Id(sdf01.getId());
|
||||||
|
String errorId = error.getSubject().getId().toString();
|
||||||
|
sDf02.setResult(errorId.substring(errorId.length() - 3));
|
||||||
|
return sDf02;
|
||||||
|
}
|
||||||
|
|
||||||
|
private SDf02 createSuccessSdf02(SDf01 sdf01) {
|
||||||
|
SDf02 sDf02 = new SDf02();
|
||||||
|
sDf02.setCurr_code(sdf01.getCurr_code());
|
||||||
|
sDf02.setAccount(sdf01.getAccount());
|
||||||
|
sDf02.setRemainder(sdf01.getRemainder());
|
||||||
|
sDf02.setDeal(sdf01.getDeal());
|
||||||
|
sDf02.setAcc_code(sdf01.getAcc_code());
|
||||||
|
sDf02.setDat(sdf01.getDat());
|
||||||
|
sDf02.setMarket(sdf01.getMarket());
|
||||||
|
sDf02.setAcc_name(sdf01.getAcc_name());
|
||||||
|
sDf02.setAcc_type(sdf01.getAcc_type());
|
||||||
|
sDf02.setSumengage(sdf01.getSumengage());
|
||||||
|
sDf02.setSumunblock(sdf01.getSumunblock());
|
||||||
|
sDf02.setFile_type(sdf01.getFile_type());
|
||||||
|
sDf02.setInSDf01Id(sdf01.getId());
|
||||||
|
sDf02.setResult("OK!");
|
||||||
|
return sDf02;
|
||||||
|
}
|
||||||
|
|
||||||
|
private AccountBalance createSuccessAccountBalance() {
|
||||||
|
return new AccountBalance();
|
||||||
|
}
|
||||||
|
|
||||||
|
private static class AccountBalance extends SpcexObjectBase {
|
||||||
|
private String message = "AccountBalance";
|
||||||
|
|
||||||
|
public String getMessage() {
|
||||||
|
return message;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setMessage(String message) {
|
||||||
|
this.message = message;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
package ru.spcex.clearing.balance.validation;
|
package ru.spcex.clearing.balance.validation;
|
||||||
|
|
||||||
import ru.clearing.classes.statics.data.account.Account;
|
import ru.clearing.classes.statics.data.account.Account;
|
||||||
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf01;
|
import ru.clearing.classes.statics.data.sdf.SDf01;
|
||||||
import ru.spcex.clearing.balance.config.ImdgValidationContext;
|
import ru.spcex.clearing.balance.config.ImdgValidationContext;
|
||||||
import ru.spcex.clearing.balance.errors.BalanceError;
|
import ru.spcex.clearing.balance.errors.BalanceError;
|
||||||
|
|
@ -11,10 +12,24 @@ import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
import ru.spcex.platform.utils.validation.IValidationRule;
|
import ru.spcex.platform.utils.validation.IValidationRule;
|
||||||
|
|
||||||
import java.time.LocalDate;
|
import java.time.LocalDate;
|
||||||
|
import java.time.format.DateTimeFormatter;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
|
|
||||||
public enum Sdf01ValidationRule implements IValidationRule<ImdgValidationContext<SDf01>> {
|
public enum Sdf01ValidationRule implements IValidationRule<ImdgValidationContext<SDf01>> {
|
||||||
|
CompanyPresent() {
|
||||||
|
@Override
|
||||||
|
public Optional<EnumMessage> validate(ImdgValidationContext<SDf01> context) {
|
||||||
|
SDf01 sdf01 = context.getValidatedObject();
|
||||||
|
Imdg<Company> companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
|
||||||
|
Company company = companyImdg.getSingleObjectByFieldValues(Map.of("tradingCode", sdf01.getDeal()));
|
||||||
|
if (company == null) {
|
||||||
|
return of( BalanceError.CompanyNotFound);
|
||||||
|
}
|
||||||
|
context.storeObject(ValidationStored.Sdf01Company, company);
|
||||||
|
return empty();
|
||||||
|
}
|
||||||
|
},
|
||||||
AccountPresent() {
|
AccountPresent() {
|
||||||
@Override
|
@Override
|
||||||
public Optional<EnumMessage> validate(ImdgValidationContext<SDf01> context) {
|
public Optional<EnumMessage> validate(ImdgValidationContext<SDf01> context) {
|
||||||
|
|
@ -44,7 +59,7 @@ public enum Sdf01ValidationRule implements IValidationRule<ImdgValidationContext
|
||||||
public Optional<EnumMessage> validate(ImdgValidationContext<SDf01> context) {
|
public Optional<EnumMessage> validate(ImdgValidationContext<SDf01> context) {
|
||||||
SDf01 sdf01 = context.getValidatedObject();
|
SDf01 sdf01 = context.getValidatedObject();
|
||||||
//fixme string format???
|
//fixme string format???
|
||||||
if (!LocalDate.now().toString().equals(sdf01.getDat())) {
|
if (!LocalDate.now().equals(LocalDate.parse(sdf01.getDat(), datFormatter))) {
|
||||||
return of(BalanceError.CurrentDateOnly);
|
return of(BalanceError.CurrentDateOnly);
|
||||||
}
|
}
|
||||||
return empty();
|
return empty();
|
||||||
|
|
@ -70,6 +85,8 @@ public enum Sdf01ValidationRule implements IValidationRule<ImdgValidationContext
|
||||||
return empty();
|
return empty();
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
private final static DateTimeFormatter datFormatter = DateTimeFormatter.ofPattern("dd.MM.yy");
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String ruleName() {
|
public String ruleName() {
|
||||||
return "Sdf01ValidationRule." + name();
|
return "Sdf01ValidationRule." + name();
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
package ru.spcex.clearing.balance.validation;
|
package ru.spcex.clearing.balance.validation;
|
||||||
|
|
||||||
public enum ValidationStored {
|
public enum ValidationStored {
|
||||||
Sdf01Account;
|
Sdf01Account, Sdf01Company;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
|
||||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
|
||||||
public enum OperationStatus implements IEnumKey {
|
public enum OperationStatus implements IEnumKey {
|
||||||
Pending("PEND");
|
Pending("PEND"), Executed("EXEC");
|
||||||
|
|
||||||
private final String key;
|
private final String key;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -107,5 +107,7 @@ public final class IMDGDistributedNames {
|
||||||
public static final String Map_KeyRate = "Map_KeyRate";
|
public static final String Map_KeyRate = "Map_KeyRate";
|
||||||
public static final String Map_MoneyMarketSecurity = "Map_MoneyMarketSecurity";
|
public static final String Map_MoneyMarketSecurity = "Map_MoneyMarketSecurity";
|
||||||
|
|
||||||
|
public static final String Map_AccountBalance = "Map_AccountBalance";
|
||||||
|
|
||||||
public static final String MAP_SEQUENCE_NAME = "MAP_SEQUENCE_NAME"; // todo вынести idGenerator отдельно и завернуть в метод, чтобы не напрямую обращаться.
|
public static final String MAP_SEQUENCE_NAME = "MAP_SEQUENCE_NAME"; // todo вынести idGenerator отдельно и завернуть в метод, чтобы не напрямую обращаться.
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -27,7 +27,8 @@ public interface Consts {
|
||||||
String USER_LOGOUT_SUCCESS = "user-logout-success";
|
String USER_LOGOUT_SUCCESS = "user-logout-success";
|
||||||
|
|
||||||
//todo
|
//todo
|
||||||
String STATEMENT_NEW = "statement-action";
|
String STATEMENT_PROCESS = "statement-process";
|
||||||
String ACCOUNT_NEW = "account-new";
|
String ACCOUNT_NEW = "account-new";
|
||||||
|
String BALANCE_ACCOUNT_NEW = "balance-account-new";
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,14 +0,0 @@
|
||||||
package ru.spcex.clearing.platform.messaging.domain.cud.account;
|
|
||||||
|
|
||||||
public class AccountNewRequest {
|
|
||||||
|
|
||||||
private String account;
|
|
||||||
|
|
||||||
public String getAccount() {
|
|
||||||
return account;
|
|
||||||
}
|
|
||||||
|
|
||||||
public void setAccount(String account) {
|
|
||||||
this.account = account;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -0,0 +1,29 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
public class AccountSdf01Request {
|
||||||
|
|
||||||
|
private Long groupingSdf01Id;
|
||||||
|
|
||||||
|
@JsonProperty("accounts")
|
||||||
|
List<AccountSdf01RequestPart> accounts;
|
||||||
|
|
||||||
|
public List<AccountSdf01RequestPart> getAccounts() {
|
||||||
|
return accounts;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setAccounts(List<AccountSdf01RequestPart> accounts) {
|
||||||
|
this.accounts = accounts;
|
||||||
|
}
|
||||||
|
|
||||||
|
public Long getGroupingSdf01Id() {
|
||||||
|
return groupingSdf01Id;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setGroupingSdf01Id(Long groupingSdf01Id) {
|
||||||
|
this.groupingSdf01Id = groupingSdf01Id;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,36 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
|
public class AccountSdf01RequestPart {
|
||||||
|
@JsonProperty
|
||||||
|
private Long sdf01Id;
|
||||||
|
@JsonProperty
|
||||||
|
private String account;
|
||||||
|
@JsonProperty
|
||||||
|
private Long companyId;
|
||||||
|
|
||||||
|
public Long getSdf01Id() {
|
||||||
|
return sdf01Id;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setSdf01Id(Long sdf01Id) {
|
||||||
|
this.sdf01Id = sdf01Id;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getAccount() {
|
||||||
|
return account;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setAccount(String account) {
|
||||||
|
this.account = account;
|
||||||
|
}
|
||||||
|
|
||||||
|
public Long getCompanyId() {
|
||||||
|
return companyId;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setCompanyId(Long companyId) {
|
||||||
|
this.companyId = companyId;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,36 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.domain.cud.balance;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
|
public class AccountSdf01ToStatementRequestPart {
|
||||||
|
@JsonProperty
|
||||||
|
private Long sdf01Id;
|
||||||
|
@JsonProperty
|
||||||
|
private Long errorCode;
|
||||||
|
@JsonProperty
|
||||||
|
private String errorText;
|
||||||
|
|
||||||
|
public Long getSdf01Id() {
|
||||||
|
return sdf01Id;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setSdf01Id(Long sdf01Id) {
|
||||||
|
this.sdf01Id = sdf01Id;
|
||||||
|
}
|
||||||
|
|
||||||
|
public Long getErrorCode() {
|
||||||
|
return errorCode;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setErrorCode(Long errorCode) {
|
||||||
|
this.errorCode = errorCode;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getErrorText() {
|
||||||
|
return errorText;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setErrorText(String errorText) {
|
||||||
|
this.errorText = errorText;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -2,15 +2,40 @@ package ru.spcex.clearing.platform.messaging.domain.cud.balance;
|
||||||
|
|
||||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
public class StatementRequest {
|
public class StatementRequest {
|
||||||
@JsonProperty
|
@JsonProperty
|
||||||
private Long sdf01Id;
|
private Long sdf01GroupId;
|
||||||
|
//-------------------
|
||||||
|
//from accountBalance creation
|
||||||
|
@JsonProperty
|
||||||
|
List<AccountSdf01ToStatementRequestPart> accountCreationResults = new ArrayList<>();
|
||||||
|
//from sDf02 creation
|
||||||
|
|
||||||
public Long getSdf01Id() {
|
// public StatementRequestType getType() {
|
||||||
return sdf01Id;
|
// if (errorCode != null || errorText != null || status != null) {
|
||||||
|
// return StatementRequestType.accountBalanceResponse;
|
||||||
|
// } else if (inOutSDfType != null) {
|
||||||
|
// return StatementRequestType.sdf02Response;
|
||||||
|
// } else return StatementRequestType.create;
|
||||||
|
// return type;
|
||||||
|
// }
|
||||||
|
|
||||||
|
public Long getSdf01GroupId() {
|
||||||
|
return sdf01GroupId;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setSdf01Id(Long sdf01Id) {
|
public void setSdf01GroupId(Long sdf01GroupId) {
|
||||||
this.sdf01Id = sdf01Id;
|
this.sdf01GroupId = sdf01GroupId;
|
||||||
|
}
|
||||||
|
|
||||||
|
public List<AccountSdf01ToStatementRequestPart> getAccountCreationResults() {
|
||||||
|
return accountCreationResults;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setAccountCreationResults(List<AccountSdf01ToStatementRequestPart> accountCreationResults) {
|
||||||
|
this.accountCreationResults = accountCreationResults;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -0,0 +1,6 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.domain.cud.balance;
|
||||||
|
|
||||||
|
public enum StatementRequestType {
|
||||||
|
create, accountBalanceResponse, sdf02Response,
|
||||||
|
batch
|
||||||
|
}
|
||||||
Loading…
Add table
Reference in a new issue