Merge branch 'CLS-742' into new

This commit is contained in:
AKurakin 2024-09-12 13:46:03 +03:00
commit 67eab3a186
9 changed files with 41 additions and 29 deletions

View file

@ -793,7 +793,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
} }
if (returnTrigger) { if (returnTrigger) {
log.debug("Remove from waiting list accountId={} (list groupId={})", accountId, inWaitingLst.getGroupId()); log.debug("Remove from waiting list accountId={} (list groupId={})", accountId, inWaitingLst.getGroupId());
waitingList.remove(inWaitingLst); waitingList.remove(accountId);
} }
if (returnTrigger) if (returnTrigger)
return inWaitingLst; return inWaitingLst;

View file

@ -132,7 +132,7 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali
clientCodeNewRequest.setCompanyId(company.getId()); clientCodeNewRequest.setCompanyId(company.getId());
clientCodeNewRequest.setCode(tkrAccount.getClientCode()); clientCodeNewRequest.setCode(tkrAccount.getClientCode());
clientCodeNewRequest.setDepoAccountId(depoAccountId); clientCodeNewRequest.setDepoAccountId(depoAccountId);
Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB); Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB.getKey());
clientCodeNewRequest.setMoneyAccountId(moneyAccount.getId()); clientCodeNewRequest.setMoneyAccountId(moneyAccount.getId());
List<Long> foreignCurrencyList = accountByCurrency.entrySet() List<Long> foreignCurrencyList = accountByCurrency.entrySet()

View file

@ -4,6 +4,8 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.DeserializationFeature; import com.fasterxml.jackson.databind.DeserializationFeature;
import com.fasterxml.jackson.databind.MapperFeature; import com.fasterxml.jackson.databind.MapperFeature;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.nio.file.Files;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
@ -12,9 +14,6 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.Resource; import org.springframework.core.io.Resource;
import ru.spcex.clearing.backendapi.meta.MetaServer; import ru.spcex.clearing.backendapi.meta.MetaServer;
import java.io.IOException;
import java.nio.file.Files;
@Configuration @Configuration
public class MetaConfiguration { public class MetaConfiguration {
@Value("file:${spring.config.location}/meta.json") @Value("file:${spring.config.location}/meta.json")

View file

@ -1,5 +1,8 @@
package ru.spcex.clearing.service.payment; package ru.spcex.clearing.service.payment;
import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
import java.math.BigDecimal; import java.math.BigDecimal;
import java.time.Instant; import java.time.Instant;
import java.util.Locale; import java.util.Locale;
@ -40,7 +43,6 @@ import ru.spcex.clearing.service.registry.RegistryManager;
import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.schedule.TradingTimeService;
import ru.spcex.clearing.service.validation.ValidationStored; import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.session.stage.impl.GatewayRequester; import ru.spcex.clearing.session.stage.impl.GatewayRequester;
import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError;
import ru.spcex.clearing.util.LocaleUtil; import ru.spcex.clearing.util.LocaleUtil;
import ru.spcex.clearing.util.security.UserRoleVerification; import ru.spcex.clearing.util.security.UserRoleVerification;
import ru.spcex.platform.enumeration.CurrencyCode; import ru.spcex.platform.enumeration.CurrencyCode;
@ -55,7 +57,6 @@ import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId; import ru.spcex.platform.utils.enumeration.IEnumId;
import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.localization.SpringPropertiesLocalizer; import ru.spcex.platform.utils.localization.SpringPropertiesLocalizer;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
import ru.spcex.platform.utils.text.TextUtil; import ru.spcex.platform.utils.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidator; import ru.spcex.platform.utils.validation.IValidator;
@ -189,8 +190,8 @@ public class PaymentInstructionOutboundService {
sDf54.setGenerationId(pmt.getId()); sDf54.setGenerationId(pmt.getId());
if (time.isTradingTime()) { if (time.isTradingTime()) {
Supplier<AssetOperationRequest> builder = () -> GatewayRequestCreator.from(pmt, Supplier<AssetOperationRequest> builder = () -> GatewayRequestCreator.from(pmt,
assets.get().a__b().getTradingCode(), assets.get().a__b().getTradingCode(),
tcr.getCode()); tcr == null ? null : tcr.getCode());
Optional<Boolean> gatewayOk = gateway.gatewayRequestAndWait(builder); Optional<Boolean> gatewayOk = gateway.gatewayRequestAndWait(builder);
if (gatewayOk.isEmpty() || gatewayOk.get().equals(Boolean.FALSE)) { if (gatewayOk.isEmpty() || gatewayOk.get().equals(Boolean.FALSE)) {
String gtwErrMsg = msgs.resolve(mapError(gatewayOk)); String gtwErrMsg = msgs.resolve(mapError(gatewayOk));

View file

@ -60,7 +60,7 @@ public class DbfImportKafkaMessenger implements InitializingBean {
} }
public void notifyUserAboutErrorParsing(Throwable error, ResultContainer resultContainer) { public void notifyUserAboutErrorParsing(Throwable error, ResultContainer resultContainer) {
if (resultContainer == null || resultContainer.getDbfFile() == null || resultContainer.getDbfFile() == null) if (resultContainer == null || resultContainer.getDbfFile() == null || resultContainer.getDbfTable() == null)
return; return;
ETable currTable = resultContainer.getDbfTable(); ETable currTable = resultContainer.getDbfTable();
String fileName = resultContainer.getDbfFile().getName(); String fileName = resultContainer.getDbfFile().getName();

View file

@ -213,8 +213,7 @@ public class TradeImporterService {
private boolean isValidTrades(STrades trades) { private boolean isValidTrades(STrades trades) {
return trades.getTradeDate() != null return trades.getTradeDate() != null
&& trades.getTradeNum() != null && trades.getTradeNum() != null
&& StringUtils.hasText(trades.getOperation()) && StringUtils.hasText(trades.getOperation());
&& trades.getTradeNum() != null;
} }
public STrades readSTrades(ResultSet resultSet) throws SQLException { public STrades readSTrades(ResultSet resultSet) throws SQLException {

View file

@ -112,7 +112,11 @@ public abstract class HazelcastServiceBase
break; break;
} }
}); });
} catch (Throwable e) { } catch (InterruptedException e) {
log.info("Hazelcast: event interrupted");
Thread.currentThread().interrupt();
}
catch (Throwable e) {
log.error("reinitializeHazelcastClient() has error: {}", ExceptionUtils.getStackTrace(e)); log.error("reinitializeHazelcastClient() has error: {}", ExceptionUtils.getStackTrace(e));
} }
} }

View file

@ -2,6 +2,17 @@ package ru.spcex.clearing.platform.messaging.service;
import com.fasterxml.jackson.databind.JavaType; import com.fasterxml.jackson.databind.JavaType;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import java.time.Duration;
import java.time.temporal.ChronoUnit;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.ConsumerRecords;
@ -24,17 +35,6 @@ import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.log.ExceptionUtils; import ru.spcex.platform.utils.log.ExceptionUtils;
import ru.spcex.platform.utils.validation.IValidator; import ru.spcex.platform.utils.validation.IValidator;
import java.time.Duration;
import java.time.temporal.ChronoUnit;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
/** /**
* утилитный класс для обработки сообщений из очереди * утилитный класс для обработки сообщений из очереди
*/ */
@ -129,6 +129,7 @@ public class QueueConsumer implements AutoCloseable {
try { try {
Thread.sleep(1000L); Thread.sleep(1000L);
} catch (InterruptedException ie) { } catch (InterruptedException ie) {
Thread.currentThread().interrupt();
log.info("Thread interrupted. {}", ExceptionUtils.getStackTrace(ie)); log.info("Thread interrupted. {}", ExceptionUtils.getStackTrace(ie));
break; break;
} }
@ -196,7 +197,10 @@ public class QueueConsumer implements AutoCloseable {
} }
send = producer.send(respRec); send = producer.send(respRec);
send.get(); send.get();
} catch (Exception e) { } catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error(ExceptionUtils.getStackTrace(e));
} catch (ExecutionException e) {
log.error(ExceptionUtils.getStackTrace(e)); log.error(ExceptionUtils.getStackTrace(e));
} }
}); });

View file

@ -1,19 +1,20 @@
package ru.spcex.clearing.platform.messaging.service; package ru.spcex.clearing.platform.messaging.service;
import com.fasterxml.jackson.databind.JavaType;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.time.Duration; import java.time.Duration;
import java.time.temporal.ChronoUnit; import java.time.temporal.ChronoUnit;
import java.util.HashMap; import java.util.HashMap;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.Future; import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function; import java.util.function.Function;
import com.fasterxml.jackson.databind.JavaType;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.ConsumerRecords;
@ -103,7 +104,8 @@ public class QueueConsumerV2 implements AutoCloseable {
log.warn("Too many error at row, {}. Sleep.", lastErrors); log.warn("Too many error at row, {}. Sleep.", lastErrors);
try { try {
Thread.sleep(1000L); Thread.sleep(1000L);
} catch (InterruptedException e){ } catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.info("Thread interrupted. {}", ExceptionUtils.getStackTrace(e)); log.info("Thread interrupted. {}", ExceptionUtils.getStackTrace(e));
break; break;
} }
@ -215,7 +217,10 @@ public class QueueConsumerV2 implements AutoCloseable {
} }
send = producer.send(respRec); send = producer.send(respRec);
send.get(); send.get();
} catch (Exception e) { } catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error(ExceptionUtils.getStackTrace(e));
} catch (ExecutionException e) {
log.error(ExceptionUtils.getStackTrace(e)); log.error(ExceptionUtils.getStackTrace(e));
} }
}); });